diff --git a/communication/src/main/java/datadog/communication/serialization/msgpack/MsgPackWriter.java b/communication/src/main/java/datadog/communication/serialization/msgpack/MsgPackWriter.java index 4fc9e8f967a..c86ce35bcac 100644 --- a/communication/src/main/java/datadog/communication/serialization/msgpack/MsgPackWriter.java +++ b/communication/src/main/java/datadog/communication/serialization/msgpack/MsgPackWriter.java @@ -101,6 +101,8 @@ public boolean format(T message, Mapper mapper) { } } buffer.reset(); + // the buffer is now empty, so drop any mapper state from the rejected message + mapper.reset(); return false; } } diff --git a/communication/src/test/java/datadog/communication/serialization/msgpack/MsgPackWriterTest.java b/communication/src/test/java/datadog/communication/serialization/msgpack/MsgPackWriterTest.java index fca302dd3d1..ecd57f010e5 100644 --- a/communication/src/test/java/datadog/communication/serialization/msgpack/MsgPackWriterTest.java +++ b/communication/src/test/java/datadog/communication/serialization/msgpack/MsgPackWriterTest.java @@ -14,6 +14,7 @@ import datadog.communication.serialization.Mapper; import datadog.communication.serialization.MessageFormatter; import datadog.communication.serialization.StreamingBuffer; +import datadog.communication.serialization.Writable; import datadog.trace.bootstrap.instrumentation.api.UTF8BytesString; import datadog.trace.util.stacktrace.StackTraceEvent; import datadog.trace.util.stacktrace.StackTraceFrame; @@ -65,6 +66,38 @@ public void testInsertAfterOverflow() { packer.format("abcdefghijklmnopqrstuvwxy", mapper), "data fits in buffer after overflow"); } + @Test + public void testMapperResetWhenOversizedMessageRejectedFromEmptyBuffer() { + CountingResetMapper mapper = new CountingResetMapper(); + MessageFormatter packer = new MsgPackWriter(newBuffer(2 + 25, (messageCount, buffer) -> {})); + assertFalse(packer.format("abcdefghijklmnopqrstuvwxyz", mapper)); + assertEquals(1, mapper.resets); + } + + @Test + public void testMapperResetWhenOversizedMessageRejectedAfterFlush() { + CountingResetMapper mapper = new CountingResetMapper(); + MessageFormatter packer = new MsgPackWriter(newBuffer(2 + 25, (messageCount, buffer) -> {})); + assertTrue(packer.format("abc", mapper)); + assertFalse(packer.format("abcdefghijklmnopqrstuvwxyz", mapper)); + // once before the retry, once after the retry is rejected + assertEquals(2, mapper.resets); + } + + private static final class CountingResetMapper implements Mapper { + int resets; + + @Override + public void map(String data, Writable writable) { + writable.writeString(data, null); + } + + @Override + public void reset() { + resets++; + } + } + @Test public void testFlushOfOverflow() { final List flushed = new ArrayList<>(); diff --git a/dd-trace-api/src/main/java/datadog/trace/api/DDTags.java b/dd-trace-api/src/main/java/datadog/trace/api/DDTags.java index b7ca19232c6..756048551c9 100644 --- a/dd-trace-api/src/main/java/datadog/trace/api/DDTags.java +++ b/dd-trace-api/src/main/java/datadog/trace/api/DDTags.java @@ -105,4 +105,5 @@ public class DDTags { public static final String PROCESS_TAGS = "_dd.tags.process"; public static final String DD_INTEGRATION = "_dd.integration"; public static final String DD_SVC_SRC = "_dd.svc_src"; + public static final String SDK_OTLP_EXPORT = "_dd.sdk.otlp_export"; } diff --git a/dd-trace-core/src/main/java/datadog/trace/common/writer/ddagent/TraceMapper.java b/dd-trace-core/src/main/java/datadog/trace/common/writer/ddagent/TraceMapper.java index 0dd3eb5e41a..e6abb09bb3c 100644 --- a/dd-trace-core/src/main/java/datadog/trace/common/writer/ddagent/TraceMapper.java +++ b/dd-trace-core/src/main/java/datadog/trace/common/writer/ddagent/TraceMapper.java @@ -12,4 +12,5 @@ public interface TraceMapper extends RemoteMapper { UTF8BytesString.create(DDSpanContext.PRIORITY_SAMPLING_KEY); static final UTF8BytesString ORIGIN_KEY = UTF8BytesString.create(DDTags.ORIGIN_KEY); static final UTF8BytesString PROCESS_TAGS_KEY = UTF8BytesString.create(DDTags.PROCESS_TAGS); + static final UTF8BytesString SDK_OTLP_EXPORT_KEY = UTF8BytesString.create(DDTags.SDK_OTLP_EXPORT); } 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 58fd278cf43..cbdcbd41433 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 @@ -101,12 +101,14 @@ public void accept(Metadata metadata) { final boolean writeSamplingPriority = firstSpanInTrace || lastSpanInTrace || metadata.topLevel(); final UTF8BytesString processTags = firstSpanInPayload ? metadata.processTags() : null; + final UTF8BytesString otlpExport = firstSpanInPayload ? metadata.otlpExportMarker() : null; int metaSize = metadata.getBaggage().size() + tags.size() + (UNSET_STATUS == metadata.getHttpStatusCode() ? 0 : 1) + (null == metadata.getOrigin() ? 0 : 1) + (null == processTags ? 0 : 1) + + (null == otlpExport ? 0 : 1) + 1; int metricsSize = (writeSamplingPriority && metadata.hasSamplingPriority() ? 1 : 0) @@ -206,6 +208,10 @@ public void accept(Metadata metadata) { writable.writeUTF8(PROCESS_TAGS_KEY); writable.writeUTF8(processTags); } + if (otlpExport != null) { + writable.writeUTF8(SDK_OTLP_EXPORT_KEY); + writable.writeUTF8(otlpExport); + } tags.forEach( writable, 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 60f221402d5..51a7f1bdb8c 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 @@ -35,6 +35,7 @@ public final class TraceMapperV0_5 implements TraceMapper { private final GrowableBuffer dictionary; private final MetaWriter metaWriter = new MetaWriter(); + private final int size; private boolean firstSpanWritten; @@ -220,6 +221,7 @@ public void accept(Metadata metadata) { final boolean writeSamplingPriority = firstSpanInTrace || lastSpanInTrace || metadata.topLevel(); final UTF8BytesString processTags = firstSpanInPayload ? metadata.processTags() : null; + final UTF8BytesString otlpExport = firstSpanInPayload ? metadata.otlpExportMarker() : null; TagMap tags = metadata.getTags(); @@ -229,6 +231,7 @@ public void accept(Metadata metadata) { + (UNSET_STATUS == metadata.getHttpStatusCode() ? 0 : 1) + (null == metadata.getOrigin() ? 0 : 1) + (null == processTags ? 0 : 1) + + (null == otlpExport ? 0 : 1) + 1; int metricsSize = (writeSamplingPriority && metadata.hasSamplingPriority() ? 1 : 0) @@ -272,6 +275,10 @@ public void accept(Metadata metadata) { writeDictionaryEncoded(writable, PROCESS_TAGS_KEY); writeDictionaryEncoded(writable, processTags); } + if (null != otlpExport) { + writeDictionaryEncoded(writable, SDK_OTLP_EXPORT_KEY); + writeDictionaryEncoded(writable, otlpExport); + } for (TagMap.EntryReader entry : tags) { if (entry.isNumber()) continue; 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 6a0ac4d7b3b..bf80e29637b 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 @@ -665,8 +665,11 @@ private ByteBuffer buildHeader() { // attributes = 10, a collection of key to value pairs common in all `chunks` CharSequence processTags = ProcessTags.getTagsForSerialization(); - Map tags = - processTags != null ? singletonMap(DDTags.PROCESS_TAGS, processTags) : emptyMap(); + Map tags = new HashMap<>(4); + tags.put(DDTags.SDK_OTLP_EXPORT, String.valueOf(cfg.isOtlpTracesExportEnabled())); + if (processTags != null) { + tags.put(DDTags.PROCESS_TAGS, processTags); + } encodeAttributes(headerWriter, 10, tags); // chunks = 11, a list of trace `chunks`, value is written by PayloadV1 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 1b211b5fae1..d5cb80b286a 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 @@ -71,6 +71,9 @@ public class DDSpanContext public static final String SPAN_SAMPLING_RULE_RATE_TAG = "_dd.span_sampling.rule_rate"; public static final String SPAN_SAMPLING_MAX_PER_SECOND_TAG = "_dd.span_sampling.max_per_second"; + private static final UTF8BytesString OTLP_EXPORT_TRUE = UTF8BytesString.create("true"); + private static final UTF8BytesString OTLP_EXPORT_FALSE = UTF8BytesString.create("false"); + private static final DDCache THREAD_NAMES = DDCaches.newFixedSizeCache(256); @@ -1416,6 +1419,7 @@ void processTagsAndBaggage( getOrigin(), longRunningVersion, ProcessTags.getTagsForSerialization(), + Config.get().isOtlpTracesExportEnabled() ? OTLP_EXPORT_TRUE : OTLP_EXPORT_FALSE, restrictedSpan.getLinks())); } } diff --git a/dd-trace-core/src/main/java/datadog/trace/core/Metadata.java b/dd-trace-core/src/main/java/datadog/trace/core/Metadata.java index d957358d73a..909e9d6dff2 100644 --- a/dd-trace-core/src/main/java/datadog/trace/core/Metadata.java +++ b/dd-trace-core/src/main/java/datadog/trace/core/Metadata.java @@ -23,6 +23,7 @@ public final class Metadata { private final CharSequence origin; private final int longRunningVersion; private final UTF8BytesString processTags; + private final UTF8BytesString otlpExportMarker; private final List spanLinks; public Metadata( @@ -37,6 +38,7 @@ public Metadata( CharSequence origin, int longRunningVersion, UTF8BytesString processTags, + UTF8BytesString otlpExportMarker, List spanLinks) { this.threadId = threadId; this.threadName = threadName; @@ -49,6 +51,7 @@ public Metadata( this.origin = origin; this.longRunningVersion = longRunningVersion; this.processTags = processTags; + this.otlpExportMarker = otlpExportMarker; this.spanLinks = spanLinks == null ? emptyList() : spanLinks; } @@ -121,6 +124,10 @@ public UTF8BytesString processTags() { return processTags; } + public UTF8BytesString otlpExportMarker() { + return otlpExportMarker; + } + public List getSpanLinks() { return spanLinks; } diff --git a/dd-trace-core/src/main/java/datadog/trace/core/otlp/common/OtlpResourceAttributes.java b/dd-trace-core/src/main/java/datadog/trace/core/otlp/common/OtlpResourceAttributes.java index 988307c6b03..6489c7ae5bf 100644 --- a/dd-trace-core/src/main/java/datadog/trace/core/otlp/common/OtlpResourceAttributes.java +++ b/dd-trace-core/src/main/java/datadog/trace/core/otlp/common/OtlpResourceAttributes.java @@ -1,6 +1,7 @@ package datadog.trace.core.otlp.common; import static datadog.communication.ddagent.TracerVersion.TRACER_VERSION; +import static datadog.trace.api.DDTags.SDK_OTLP_EXPORT; import static java.util.Arrays.asList; import datadog.trace.api.Config; @@ -23,6 +24,8 @@ private OtlpResourceAttributes() {} /** Marks that the Agent should not recompute trace metrics from the exported spans. */ private static final String STATS_COMPUTED_KEY = "_dd.stats_computed"; + private static final String SDK_SEMANTICS_KEY = "datadog.sdk.semantics"; + private static final Set IGNORED_GLOBAL_TAGS = new HashSet<>( asList( @@ -34,7 +37,9 @@ private OtlpResourceAttributes() {} "service.version", "telemetry.sdk.name", "telemetry.sdk.version", - "telemetry.sdk.language")); + "telemetry.sdk.language", + SDK_SEMANTICS_KEY, + SDK_OTLP_EXPORT)); /** * {@code value} is a {@link String}, except {@code datadog.process_tags}: a {@code List}. @@ -81,10 +86,15 @@ static void visitResourceAttributes( /** * Builds the extra resource attributes for the OTLP trace export: the {@code _dd.stats_computed} * marker when the SDK is computing OTLP span metrics, so a downstream Agent does not recompute - * them from the exported spans. + * them from the exported spans; {@code datadog.sdk.semantics}, recording whether the SDK applied + * Datadog or OTel semantics; and {@code _dd.sdk.otlp_export}, which is always {@code "true"} here + * because reaching this encoder means the payload is leaving over OTLP. */ static Map traceResourceAttributes(Config config) { Map attributes = new LinkedHashMap<>(); + attributes.put( + "datadog.sdk.semantics", config.isTraceOtelSemanticsEnabled() ? "otel" : "datadog"); + attributes.put(SDK_OTLP_EXPORT, "true"); if (config.isOtelTracesSpanMetricsEnabled()) { attributes.put(STATS_COMPUTED_KEY, "true"); } diff --git a/dd-trace-core/src/test/java/datadog/trace/common/writer/DDAgentApiTest.java b/dd-trace-core/src/test/java/datadog/trace/common/writer/DDAgentApiTest.java index 4e0abf5fbab..66b50ff1dde 100644 --- a/dd-trace-core/src/test/java/datadog/trace/common/writer/DDAgentApiTest.java +++ b/dd-trace-core/src/test/java/datadog/trace/common/writer/DDAgentApiTest.java @@ -1,5 +1,6 @@ package datadog.trace.common.writer; +import static datadog.trace.api.DDTags.SDK_OTLP_EXPORT; import static datadog.trace.api.ProtocolVersion.V0_5; import static java.util.Collections.emptyList; import static java.util.Collections.emptyMap; @@ -252,6 +253,8 @@ void testContentIsSentAsMsgpackServiceSpan() throws IOException { && ProcessTags.getTagsForSerialization() != null) { meta.put("_dd.tags.process", ProcessTags.getTagsForSerialization().toString()); } + // payload-scoped marker, written on the first span of the first non-empty chunk + meta.put(SDK_OTLP_EXPORT, "false"); Map metrics = new TreeMap<>(); metrics.put(DDSpanContext.PRIORITY_SAMPLING_KEY, 1); metrics.put(InstrumentationTags.DD_TOP_LEVEL.toString(), 1); @@ -328,6 +331,8 @@ void testContentIsSentAsMsgpackResourceSpan() throws IOException { && ProcessTags.getTagsForSerialization() != null) { meta.put("_dd.tags.process", ProcessTags.getTagsForSerialization().toString()); } + // payload-scoped marker, written on the first span of the first non-empty chunk + meta.put(SDK_OTLP_EXPORT, "false"); Map metrics = new TreeMap<>(); metrics.put(DDSpanContext.PRIORITY_SAMPLING_KEY, 1); metrics.put(InstrumentationTags.DD_TOP_LEVEL.toString(), 1); diff --git a/dd-trace-core/src/test/java/datadog/trace/common/writer/DDAgentWriterCombinedTest.java b/dd-trace-core/src/test/java/datadog/trace/common/writer/DDAgentWriterCombinedTest.java index ae9741fe827..3f1bc9d1324 100644 --- a/dd-trace-core/src/test/java/datadog/trace/common/writer/DDAgentWriterCombinedTest.java +++ b/dd-trace-core/src/test/java/datadog/trace/common/writer/DDAgentWriterCombinedTest.java @@ -269,8 +269,11 @@ void testDefaultBufferSizeFor(String agentVersion) { TraceMapper mapper = agentVersion.equals("v0.5/traces") ? new TraceMapperV0_5() : new TraceMapperV0_4(); - int traceSize = calculateSize(minimalTrace, mapper); - int maxedPayloadTraceCount = (mapper.messageBufferSize() / traceSize); + // The first trace of a payload is larger: it carries the payload-scoped _dd.sdk.otlp_export + // marker on its first span. Size both cases so the overflow point is exact. + int firstTraceSize = calculateSize(minimalTrace, mapper, true); + int traceSize = calculateSize(minimalTrace, mapper, false); + int maxedPayloadTraceCount = 1 + (mapper.messageBufferSize() - firstTraceSize) / traceSize; when(discovery.getTraceEndpoint()).thenReturn(agentVersion); when(api.sendSerializedTraces( @@ -824,13 +827,25 @@ void statsdCommFailure() throws Exception { healthMetrics.close(); } - static int calculateSize(List trace, TraceMapper mapper) { + /** + * Serialized size of {@code trace}, either as the first trace of a payload (which carries the + * payload-scoped markers on its first span) or as any later trace. Uses a throwaway mapper of the + * same kind so the caller's mapper keeps its state. + */ + static int calculateSize(List trace, TraceMapper mapper, boolean firstInPayload) { AtomicInteger size = new AtomicInteger(); MsgPackWriter packer = new MsgPackWriter( new FlushingBuffer( 1024, (messageCount, buffer) -> size.set(buffer.limit() - buffer.position()))); - packer.format(trace, mapper); + TraceMapper sizingMapper = + mapper instanceof TraceMapperV0_5 ? new TraceMapperV0_5() : new TraceMapperV0_4(); + if (!firstInPayload) { + // burn the payload-scoped markers on a throwaway trace + packer.format(trace, sizingMapper); + packer.flush(); + } + packer.format(trace, sizingMapper); packer.flush(); return size.get(); } diff --git a/dd-trace-core/src/test/java/datadog/trace/common/writer/FileBasedPayloadDispatcherTest.java b/dd-trace-core/src/test/java/datadog/trace/common/writer/FileBasedPayloadDispatcherTest.java index 25c26bcd10a..2eb5e86f34b 100644 --- a/dd-trace-core/src/test/java/datadog/trace/common/writer/FileBasedPayloadDispatcherTest.java +++ b/dd-trace-core/src/test/java/datadog/trace/common/writer/FileBasedPayloadDispatcherTest.java @@ -313,6 +313,7 @@ private static CoreSpan mockSpan(CharSequence type, Map tags) null, 0, null, + null, null); doAnswer( inv -> { diff --git a/dd-trace-core/src/test/java/datadog/trace/common/writer/TraceGenerator.java b/dd-trace-core/src/test/java/datadog/trace/common/writer/TraceGenerator.java index e777235fef3..93d6130df92 100644 --- a/dd-trace-core/src/test/java/datadog/trace/common/writer/TraceGenerator.java +++ b/dd-trace-core/src/test/java/datadog/trace/common/writer/TraceGenerator.java @@ -3,6 +3,7 @@ import static datadog.trace.api.sampling.PrioritySampling.UNSET; import static java.util.Collections.emptyList; +import datadog.trace.api.Config; import datadog.trace.api.DDSpanId; import datadog.trace.api.DDTags; import datadog.trace.api.DDTraceId; @@ -242,6 +243,7 @@ public PojoSpan( origin, 0, ProcessTags.getTagsForSerialization(), + UTF8BytesString.create(String.valueOf(Config.get().isOtlpTracesExportEnabled())), spanLinks); } diff --git a/dd-trace-core/src/test/java/datadog/trace/common/writer/ddagent/TraceMapperV04PayloadTest.java b/dd-trace-core/src/test/java/datadog/trace/common/writer/ddagent/TraceMapperV04PayloadTest.java index 15483b94c95..6663c7c9501 100644 --- a/dd-trace-core/src/test/java/datadog/trace/common/writer/ddagent/TraceMapperV04PayloadTest.java +++ b/dd-trace-core/src/test/java/datadog/trace/common/writer/ddagent/TraceMapperV04PayloadTest.java @@ -1,6 +1,9 @@ package datadog.trace.common.writer.ddagent; +import static datadog.trace.api.config.TracerConfig.WRITER_TYPE; import static datadog.trace.bootstrap.instrumentation.api.InstrumentationTags.DD_MEASURED; +import static datadog.trace.bootstrap.instrumentation.api.WriterConstants.MULTI_WRITER_TYPE; +import static datadog.trace.bootstrap.instrumentation.api.WriterConstants.OTLP_WRITER_TYPE; import static datadog.trace.common.writer.TraceGenerator.generateRandomTraces; import static datadog.trace.common.writer.ddagent.PayloadVerifiers.assertEqualsWithNullAsEmpty; import static datadog.trace.common.writer.ddagent.PayloadVerifiers.unpackNumber; @@ -23,10 +26,12 @@ import datadog.trace.common.writer.Payload; import datadog.trace.common.writer.TraceGenerator.PojoSpan; import datadog.trace.core.DDSpanContext; +import datadog.trace.test.junit.utils.config.WithConfig; import datadog.trace.test.junit.utils.config.WithConfigExtension; import java.io.IOException; import java.nio.ByteBuffer; import java.util.ArrayList; +import java.util.Arrays; import java.util.Collections; import java.util.HashMap; import java.util.List; @@ -63,9 +68,6 @@ void tracesWrittenCorrectly(int bufferSize, int traceCount, boolean lowCardinali if (!packer.format(trace, traceMapper)) { verifier.skipLargeTrace(); tracesFitInBuffer = false; - // in the real like the mapper is always reset each trace. - // here we need to force it when we fail since the buffer will be reset as well - traceMapper.reset(); } } packer.flush(); @@ -101,24 +103,7 @@ private static Stream tracesWrittenCorrectlyArguments() { @MethodSource("fullSixtyFourBitTraceAndSpanIdentifiersArguments") void fullSixtyFourBitTraceAndSpanIdentifiers( String scenario, DDTraceId traceId, long spanId, long parentId) { - PojoSpan span = - new PojoSpan( - "service", - "operation", - "resource", - traceId, - spanId, - parentId, - 123L, - 456L, - 0, - Collections.emptyMap(), - Collections.emptyMap(), - "type", - false, - 0, - 0, - "origin"); + PojoSpan span = plainSpan(traceId, spanId, parentId); List> traces = Collections.singletonList(Collections.singletonList(span)); TraceMapperV0_4 traceMapper = new TraceMapperV0_4(); PayloadVerifier verifier = new PayloadVerifier(traces, traceMapper); @@ -139,24 +124,7 @@ private static Stream fullSixtyFourBitTraceAndSpanIdentifiersArgument @Test void metaStructSupport() { - PojoSpan span = - new PojoSpan( - "service", - "operation", - "resource", - DDTraceId.ONE, - 1L, - -1L, - 123L, - 456L, - 0, - Collections.emptyMap(), - Collections.emptyMap(), - "type", - false, - 0, - 0, - "origin"); + PojoSpan span = plainSpan(1L); List> stack = new ArrayList<>(); for (StackTraceElement element : Thread.currentThread().getStackTrace()) { Map frame = new HashMap<>(); @@ -200,24 +168,7 @@ void processTagsSerialization() { assertNotNull(ProcessTags.getTagsForSerialization()); List spans = new ArrayList<>(); for (long spanId = 1; spanId <= 2; ++spanId) { - spans.add( - new PojoSpan( - "service", - "operation", - "resource", - DDTraceId.ONE, - spanId, - -1L, - 123L, - 456L, - 0, - Collections.emptyMap(), - Collections.emptyMap(), - "type", - false, - 0, - 0, - "origin")); + spans.add(plainSpan(spanId)); } List> traces = Collections.singletonList(spans); @@ -231,6 +182,75 @@ void processTagsSerialization() { verifier.verifyTracesConsumed(); } + /** + * v0.4 has no payload-level field, so the export-mode marker rides on the first span of the first + * non-empty chunk and the Agent hoists it onto {@code TracerPayload.tags}. + */ + @Test + void otlpExportMarkerOnlyOnFirstSpanOfFirstNonEmptyChunk() { + List> traces = + Arrays.asList( + Collections.emptyList(), + Arrays.asList(plainSpan(1), plainSpan(2)), + Collections.singletonList(plainSpan(3))); + TraceMapperV0_4 traceMapper = new TraceMapperV0_4(); + PayloadVerifier verifier = new PayloadVerifier(traces, traceMapper); + MsgPackWriter packer = new MsgPackWriter(new FlushingBuffer(200 << 10, verifier)); + + for (List trace : traces) { + assertTrue(packer.format(trace, traceMapper)); + } + packer.flush(); + + verifier.verifyTracesConsumed(); + // The verifier already asserts exactly one marker per payload, that it sits on span 0 and that + // it holds the expected value; pin down *which chunk* it landed on: the first NON-EMPTY one, + // not the leading empty one and not a later one. + assertEquals(1, verifier.otlpExportTraceIndex()); + } + + @Test + @WithConfig( + key = WRITER_TYPE, + value = MULTI_WRITER_TYPE + ":" + OTLP_WRITER_TYPE + ",DDAgentWriter") + void otlpExportMarkerIsTrueWhenAlsoExportingOverOtlp() { + List> traces = + Collections.singletonList(Collections.singletonList(plainSpan(1))); + TraceMapperV0_4 traceMapper = new TraceMapperV0_4(); + PayloadVerifier verifier = new PayloadVerifier(traces, traceMapper).expectOtlpExport("true"); + MsgPackWriter packer = new MsgPackWriter(new FlushingBuffer(200 << 10, verifier)); + + packer.format(traces.get(0), traceMapper); + packer.flush(); + + verifier.verifyTracesConsumed(); + assertEquals(0, verifier.otlpExportTraceIndex()); + } + + private static PojoSpan plainSpan(long spanId) { + return plainSpan(DDTraceId.ONE, spanId, -1L); + } + + private static PojoSpan plainSpan(DDTraceId traceId, long spanId, long parentId) { + return new PojoSpan( + "service", + "operation", + "resource", + traceId, + spanId, + parentId, + 123L, + 456L, + 0, + Collections.emptyMap(), + Collections.emptyMap(), + "type", + false, + 0, + 0, + "origin"); + } + private static final class PayloadVerifier implements ByteBufferConsumer { private final List> expectedTraces; @@ -241,6 +261,16 @@ private static final class PayloadVerifier implements ByteBufferConsumer { private int position = 0; + /** Expected value of the payload-scoped {@code _dd.sdk.otlp_export} marker. */ + private String expectedOtlpExport = "false"; + + /** + * Payload-spanning index of the chunk the marker was last seen on, or {@code -1} if it was + * never seen. The span index within that chunk is not tracked: {@link #accept} already asserts + * the marker only ever rides span 0. + */ + private int otlpExportTraceIndex = -1; + private PayloadVerifier(List> traces, TraceMapperV0_4 mapper) { this(traces, mapper, null); } @@ -254,6 +284,16 @@ private PayloadVerifier( this.metaStructVerifier = metaStructVerifier; } + /** Sets the expected {@code _dd.sdk.otlp_export} value (defaults to {@code "false"}). */ + PayloadVerifier expectOtlpExport(String value) { + this.expectedOtlpExport = value; + return this; + } + + int otlpExportTraceIndex() { + return otlpExportTraceIndex; + } + void skipLargeTrace() { ++position; } @@ -264,6 +304,7 @@ public void accept(int messageCount, ByteBuffer buffer) { return; } int processTagsCount = 0; + int otlpExportCount = 0; try { Payload payload = mapper.newPayload().withBody(messageCount, buffer); payload.writeTo(channel); @@ -358,6 +399,14 @@ public void accept(int messageCount, ByteBuffer buffer) { assertEquals(0, k); assertEquals(ProcessTags.getTagsForSerialization().toString(), entry.getValue()); processTagsCount++; + } else if (DDTags.SDK_OTLP_EXPORT.equals(entry.getKey())) { + // Payload-scoped: only the first span of the first non-empty chunk carries it. + otlpExportCount++; + assertEquals(0, k); + assertEquals(expectedOtlpExport, entry.getValue()); + // `position` was post-incremented when this trace was picked up, so the current + // trace's payload-spanning index is `position - 1`. + otlpExportTraceIndex = position - 1; } else { Object tag = expectedSpan.getTag(entry.getKey()); if (tag != null) { @@ -389,6 +438,8 @@ public void accept(int messageCount, ByteBuffer buffer) { channel.resetForWriting(); assertEquals( Config.get().isExperimentalPropagateProcessTagsEnabled() ? 1 : 0, processTagsCount); + // exactly one _dd.sdk.otlp_export per payload, never per span + assertEquals(1, otlpExportCount); } } diff --git a/dd-trace-core/src/test/java/datadog/trace/common/writer/ddagent/TraceMapperV05PayloadTest.java b/dd-trace-core/src/test/java/datadog/trace/common/writer/ddagent/TraceMapperV05PayloadTest.java index 6dedd2e9ef6..aada5c9bb7b 100644 --- a/dd-trace-core/src/test/java/datadog/trace/common/writer/ddagent/TraceMapperV05PayloadTest.java +++ b/dd-trace-core/src/test/java/datadog/trace/common/writer/ddagent/TraceMapperV05PayloadTest.java @@ -1,7 +1,10 @@ package datadog.trace.common.writer.ddagent; import static datadog.trace.api.config.GeneralConfig.EXPERIMENTAL_PROPAGATE_PROCESS_TAGS_ENABLED; +import static datadog.trace.api.config.TracerConfig.WRITER_TYPE; import static datadog.trace.bootstrap.instrumentation.api.InstrumentationTags.DD_MEASURED; +import static datadog.trace.bootstrap.instrumentation.api.WriterConstants.MULTI_WRITER_TYPE; +import static datadog.trace.bootstrap.instrumentation.api.WriterConstants.OTLP_WRITER_TYPE; import static datadog.trace.common.writer.TraceGenerator.generateRandomTraces; import static datadog.trace.common.writer.ddagent.PayloadVerifiers.assertEqualsWithNullAsEmpty; import static datadog.trace.common.writer.ddagent.PayloadVerifiers.unpackNumber; @@ -30,6 +33,7 @@ import java.io.IOException; import java.nio.ByteBuffer; import java.util.ArrayList; +import java.util.Arrays; import java.util.Collections; import java.util.HashMap; import java.util.List; @@ -95,10 +99,14 @@ void bodyOverflowCausesAFlush() { PrioritySampling.UNSET, 0, null)); - int traceSize = calculateSize(repeatedTrace); + // The first trace of a payload is larger: it carries the payload-scoped _dd.sdk.otlp_export + // marker on its first span. Size both cases so the overflow point is exact. + int firstTraceSize = calculateSize(repeatedTrace, true); + int traceSize = calculateSize(repeatedTrace, false); // 30KB body int bufferSize = 30 << 10; - int tracesRequiredToOverflowBody = (int) Math.ceil((double) bufferSize / traceSize) + 1; + int tracesThatFitInBody = 1 + (bufferSize - firstTraceSize) / traceSize; + int tracesRequiredToOverflowBody = tracesThatFitInBody + 1; List> traces = new ArrayList<>(tracesRequiredToOverflowBody); for (int i = 0; i < tracesRequiredToOverflowBody; ++i) { traces.add(repeatedTrace); @@ -174,24 +182,7 @@ void processTagsSerialization() { assertNotNull(ProcessTags.getTagsForSerialization()); List spans = new ArrayList<>(); for (long spanId = 1; spanId <= 2; ++spanId) { - spans.add( - new PojoSpan( - "service", - "operation", - "resource", - DDTraceId.ONE, - spanId, - -1L, - 123L, - 456L, - 0, - Collections.emptyMap(), - Collections.emptyMap(), - "type", - false, - 0, - 0, - "origin")); + spans.add(plainSpan(spanId)); } List> traces = Collections.singletonList(spans); @@ -205,6 +196,71 @@ void processTagsSerialization() { verifier.verifyTracesConsumed(); } + /** + * v0.5 has no payload-level field, so the export-mode marker rides on the first span of the first + * non-empty chunk and the Agent hoists it onto {@code TracerPayload.tags}. + */ + @Test + void otlpExportMarkerOnlyOnFirstSpanOfFirstNonEmptyChunk() { + List> traces = + Arrays.asList( + Collections.emptyList(), + Arrays.asList(plainSpan(1), plainSpan(2)), + Collections.singletonList(plainSpan(3))); + TraceMapperV0_5 traceMapper = new TraceMapperV0_5(); + PayloadVerifier verifier = new PayloadVerifier(traces, traceMapper); + MsgPackWriter packer = new MsgPackWriter(new FlushingBuffer(200 << 10, verifier)); + + for (List trace : traces) { + assertTrue(packer.format(trace, traceMapper)); + } + packer.flush(); + + verifier.verifyTracesConsumed(); + // The verifier already asserts exactly one marker per payload, that it sits on span 0 and that + // it holds the expected value; pin down *which chunk* it landed on: the first NON-EMPTY one, + // not the leading empty one and not a later one. + assertEquals(1, verifier.otlpExportTraceIndex()); + } + + @Test + @WithConfig( + key = WRITER_TYPE, + value = MULTI_WRITER_TYPE + ":" + OTLP_WRITER_TYPE + ",DDAgentWriter") + void otlpExportMarkerIsTrueWhenAlsoExportingOverOtlp() { + List> traces = + Collections.singletonList(Collections.singletonList(plainSpan(1))); + TraceMapperV0_5 traceMapper = new TraceMapperV0_5(); + PayloadVerifier verifier = new PayloadVerifier(traces, traceMapper).expectOtlpExport("true"); + MsgPackWriter packer = new MsgPackWriter(new FlushingBuffer(200 << 10, verifier)); + + packer.format(traces.get(0), traceMapper); + packer.flush(); + + verifier.verifyTracesConsumed(); + assertEquals(0, verifier.otlpExportTraceIndex()); + } + + private static PojoSpan plainSpan(long spanId) { + return new PojoSpan( + "service", + "operation", + "resource", + DDTraceId.ONE, + spanId, + -1L, + 123L, + 456L, + 0, + Collections.emptyMap(), + Collections.emptyMap(), + "type", + false, + 0, + 0, + "origin"); + } + private static final class PayloadVerifier implements ByteBufferConsumer { private final List> expectedTraces; @@ -213,6 +269,16 @@ private static final class PayloadVerifier implements ByteBufferConsumer { private int position = 0; + /** Expected value of the payload-scoped {@code _dd.sdk.otlp_export} marker. */ + private String expectedOtlpExport = "false"; + + /** + * Payload-spanning index of the chunk the marker was last seen on, or {@code -1} if it was + * never seen. The span index within that chunk is not tracked: {@link #accept} already asserts + * the marker only ever rides span 0. + */ + private int otlpExportTraceIndex = -1; + private PayloadVerifier(List> traces, TraceMapperV0_5 mapper) { this(traces, mapper, 200 << 10); } @@ -223,13 +289,27 @@ private PayloadVerifier(List> traces, TraceMapperV0_5 mapper, int this.channel = new PayloadVerifiers.CapturingChannel(size); } + /** Sets the expected {@code _dd.sdk.otlp_export} value (defaults to {@code "false"}). */ + PayloadVerifier expectOtlpExport(String value) { + this.expectedOtlpExport = value; + return this; + } + + int otlpExportTraceIndex() { + return otlpExportTraceIndex; + } + void skipLargeTrace() { ++position; } @Override public void accept(int messageCount, ByteBuffer buffer) { + if (expectedTraces.isEmpty() && messageCount == 0) { + return; + } int processTagsCount = 0; + int otlpExportCount = 0; try { Payload payload = mapper.newPayload().withBody(messageCount, buffer); payload.writeTo(channel); @@ -283,6 +363,14 @@ public void accept(int messageCount, ByteBuffer buffer) { assertTrue(Config.get().isExperimentalPropagateProcessTagsEnabled()); assertEquals(0, k); assertEquals(ProcessTags.getTagsForSerialization().toString(), entry.getValue()); + } else if (DDTags.SDK_OTLP_EXPORT.equals(entry.getKey())) { + // Payload-scoped: only the first span of the first non-empty chunk carries it. + otlpExportCount++; + assertEquals(0, k); + assertEquals(expectedOtlpExport, entry.getValue()); + // `position` was post-incremented when this trace was picked up, so the current + // trace's payload-spanning index is `position - 1`. + otlpExportTraceIndex = position - 1; } else { Object tag = expectedSpan.getTag(entry.getKey()); if (tag != null) { @@ -334,10 +422,14 @@ public void accept(int messageCount, ByteBuffer buffer) { } catch (IOException e) { fail(e.getMessage()); } finally { - assertEquals( - Config.get().isExperimentalPropagateProcessTagsEnabled() ? 1 : 0, processTagsCount); + // Reset before asserting: a failing assertion here must not leave the mapper and channel + // dirty for the next payload of a multi-payload test. mapper.reset(); channel.resetForWriting(); + assertEquals( + Config.get().isExperimentalPropagateProcessTagsEnabled() ? 1 : 0, processTagsCount); + // exactly one _dd.sdk.otlp_export per payload, never per span + assertEquals(1, otlpExportCount); } } @@ -346,7 +438,11 @@ void verifyTracesConsumed() { } } - private static int calculateSize(List trace) { + /** + * Serialized size of {@code trace}, either as the first trace of a payload (which carries the + * payload-scoped markers on its first span) or as any later trace. + */ + private static int calculateSize(List trace, boolean firstInPayload) { AtomicInteger size = new AtomicInteger(); MsgPackWriter packer = new MsgPackWriter( @@ -358,7 +454,13 @@ public void accept(int messageCount, ByteBuffer buffer) { size.set(buffer.limit() - buffer.position()); } })); - packer.format(trace, new TraceMapperV0_5(1024)); + TraceMapperV0_5 mapper = new TraceMapperV0_5(1024); + if (!firstInPayload) { + // burn the payload-scoped markers on a throwaway trace + packer.format(trace, mapper); + packer.flush(); + } + packer.format(trace, mapper); packer.flush(); return size.get(); } diff --git a/dd-trace-core/src/test/java/datadog/trace/common/writer/ddagent/TraceMapperV1PayloadTest.java b/dd-trace-core/src/test/java/datadog/trace/common/writer/ddagent/TraceMapperV1PayloadTest.java index 1eacd903b96..0da528e2eb4 100644 --- a/dd-trace-core/src/test/java/datadog/trace/common/writer/ddagent/TraceMapperV1PayloadTest.java +++ b/dd-trace-core/src/test/java/datadog/trace/common/writer/ddagent/TraceMapperV1PayloadTest.java @@ -1,9 +1,11 @@ package datadog.trace.common.writer.ddagent; import static datadog.trace.api.DDTags.PROCESS_TAGS; +import static datadog.trace.api.DDTags.SDK_OTLP_EXPORT; import static datadog.trace.api.DDTags.SPAN_EVENTS; import static datadog.trace.api.DDTags.THREAD_ID; import static datadog.trace.api.DDTags.THREAD_NAME; +import static datadog.trace.api.config.TracerConfig.WRITER_TYPE; import static datadog.trace.api.sampling.PrioritySampling.SAMPLER_KEEP; import static datadog.trace.bootstrap.instrumentation.api.Tags.COMPONENT; import static datadog.trace.bootstrap.instrumentation.api.Tags.ENV; @@ -12,6 +14,8 @@ import static datadog.trace.bootstrap.instrumentation.api.Tags.SPAN_KIND; import static datadog.trace.bootstrap.instrumentation.api.Tags.SPAN_KIND_CLIENT; import static datadog.trace.bootstrap.instrumentation.api.Tags.VERSION; +import static datadog.trace.bootstrap.instrumentation.api.WriterConstants.MULTI_WRITER_TYPE; +import static datadog.trace.bootstrap.instrumentation.api.WriterConstants.OTLP_WRITER_TYPE; import static datadog.trace.common.writer.TraceGenerator.generateRandomTraces; import static datadog.trace.common.writer.ddagent.PayloadVerifiers.assertEqualsWithNullAsEmpty; import static datadog.trace.common.writer.ddagent.V1PayloadReader.newStringTable; @@ -19,6 +23,7 @@ import static datadog.trace.common.writer.ddagent.V1PayloadReader.readBinary; import static datadog.trace.common.writer.ddagent.V1PayloadReader.readFirstChunk; import static datadog.trace.common.writer.ddagent.V1PayloadReader.readFirstSpan; +import static datadog.trace.common.writer.ddagent.V1PayloadReader.readPayloadAttributes; import static datadog.trace.common.writer.ddagent.V1PayloadReader.readStreamingString; import static datadog.trace.common.writer.ddagent.V1PayloadReader.skipChunkField; import static datadog.trace.common.writer.ddagent.V1PayloadReader.skipPayloadField; @@ -58,6 +63,7 @@ import datadog.trace.common.writer.ddagent.V1PayloadReader.V1SpanEvent; import datadog.trace.common.writer.ddagent.V1PayloadReader.V1SpanLink; import datadog.trace.core.MetadataConsumer; +import datadog.trace.test.junit.utils.config.WithConfig; import datadog.trace.test.junit.utils.config.WithConfigExtension; import java.io.ByteArrayOutputStream; import java.io.IOException; @@ -149,7 +155,6 @@ void tracesWrittenCorrectly( if (!packer.format(trace, traceMapper)) { verifier.skipLargeTrace(); tracesFitInBuffer = false; - traceMapper.reset(); } } packer.flush(); @@ -235,16 +240,33 @@ void payloadContainsExpectedHeaderAndChunkFields() throws IOException { assertEquals(EXPECTED_PAYLOAD_FIELD_IDS, payloadFieldsSeen); assertEquals(1, chunkCount); assertNotNull(payloadAttributes); + // The export-mode marker is payload-scoped and always written; process tags are conditional. + assertEquals("false", payloadAttributes.get(SDK_OTLP_EXPORT)); CharSequence processTags = ProcessTags.getTagsForSerialization(); if (processTags == null) { - assertEquals(0, payloadAttributes.size()); - } else { assertEquals(1, payloadAttributes.size()); + } else { + assertEquals(2, payloadAttributes.size()); assertEquals(processTags.toString(), payloadAttributes.get(PROCESS_TAGS)); } } } + /** + * v1 has a payload-level attribute map, so the export-mode marker lives there rather than on the + * first span the way v0.4/v0.5 must place it. The {@code "false"} default is covered by {@link + * #payloadContainsExpectedHeaderAndChunkFields()}, which also pins the attribute-map size. + */ + @Test + @WithConfig( + key = WRITER_TYPE, + value = MULTI_WRITER_TYPE + ":" + OTLP_WRITER_TYPE + ",DDAgentWriter") + void otlpExportMarkerIsTrueWhenAlsoExportingOverOtlp() throws IOException { + Map attributes = readPayloadAttributes(serializeV1Payload(span(emptyMap()))); + + assertEquals("true", attributes.get(SDK_OTLP_EXPORT)); + } + // expectedSamplingMechanism 0 is SamplingMechanism.DEFAULT. @TableTest({ "scenario | decisionMakerTag | expectedSamplingMechanism", 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 4dfdaf4bbd7..82ef0b2df65 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 @@ -150,6 +150,21 @@ public static V1Chunk readFirstChunk(byte[] encoded) throws IOException { throw new AssertionError("Could not find first chunk in v1 payload"); } + /** Decodes the payload-level attribute map (field 10) of an encoded V1 payload. */ + public static Map readPayloadAttributes(byte[] encoded) throws IOException { + MessageUnpacker unpacker = MessagePack.newDefaultUnpacker(new ArrayBufferInput(encoded)); + List stringTable = newStringTable(); + int payloadFieldCount = unpacker.unpackMapHeader(); + for (int i = 0; i < payloadFieldCount; i++) { + int payloadFieldId = unpacker.unpackInt(); + if (payloadFieldId == PayloadField.ATTRIBUTES) { + return readAttributes(unpacker, stringTable); + } + skipPayloadField(unpacker, payloadFieldId, stringTable); + } + throw new AssertionError("Could not find payload attributes in v1 payload"); + } + /** Creates a string table seeded with the empty string at index 0, as the writer expects. */ public static List newStringTable() { List stringTable = new ArrayList<>(); 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 839abe61a39..cf08059ea9f 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,5 +1,6 @@ package datadog.trace.core; +import static datadog.trace.api.DDTags.SDK_OTLP_EXPORT; import static datadog.trace.api.DDTags.SPAN_EVENTS; import static datadog.trace.api.DDTags.SPAN_LINKS; import static datadog.trace.api.TracePropagationStyle.DATADOG; @@ -245,7 +246,7 @@ void serializeTraceWithBaggageAndTagsCorrectlyV04( for (int j = 0; j < packedSize; j++) { String k = unpacker.unpackString(); String v = unpacker.unpackString(); - if (!"thread.name".equals(k) && !"thread.id".equals(k)) { + if (!isWrittenByMapper(k)) { unpackedMeta.put(k, v); } } @@ -324,7 +325,7 @@ void serializeTraceWithBaggageAndTagsCorrectlyV05( for (int j = 0; j < packedSize; j++) { String k = dictionary[unpacker.unpackInt()]; String v = dictionary[unpacker.unpackInt()]; - if (!"thread.name".equals(k) && !"thread.id".equals(k)) { + if (!isWrittenByMapper(k)) { unpackedMeta.put(k, v); } } @@ -562,7 +563,7 @@ void serializeTraceWithFlatMapTagV04() throws Exception { for (int j = 0; j < packedSize; j++) { String k = unpacker.unpackString(); String v = unpacker.unpackString(); - if (!"thread.name".equals(k) && !"thread.id".equals(k)) { + if (!isWrittenByMapper(k)) { unpackedMeta.put(k, v); } } @@ -632,7 +633,7 @@ void serializeTraceWithFlatMapTagV05() throws Exception { for (int j = 0; j < packedSize; j++) { String k = dictionary[unpacker.unpackInt()]; String v = dictionary[unpacker.unpackInt()]; - if (!"thread.name".equals(k) && !"thread.id".equals(k)) { + if (!isWrittenByMapper(k)) { unpackedMeta.put(k, v); } } @@ -640,6 +641,14 @@ void serializeTraceWithFlatMapTagV05() throws Exception { tracer.close(); } + /** + * thread.* and the payload-scoped _dd.sdk.otlp_export marker are written by the mapper, not by + * the span under test (see TraceMapperV04/V05PayloadTest). + */ + private static boolean isWrittenByMapper(String key) { + return "thread.name".equals(key) || "thread.id".equals(key) || SDK_OTLP_EXPORT.equals(key); + } + private static class CaptureBuffer implements ByteBufferConsumer { private byte[] bytes; int messageCount; diff --git a/dd-trace-core/src/test/java/datadog/trace/core/MetadataTest.java b/dd-trace-core/src/test/java/datadog/trace/core/MetadataTest.java index 2815cac6730..eab3afd2812 100644 --- a/dd-trace-core/src/test/java/datadog/trace/core/MetadataTest.java +++ b/dd-trace-core/src/test/java/datadog/trace/core/MetadataTest.java @@ -62,6 +62,7 @@ private static Metadata metadataWithStatus(int status) { null, 0, null, + null, emptyList()); } } diff --git a/dd-trace-core/src/test/java/datadog/trace/core/otlp/common/OtlpResourceJsonTest.java b/dd-trace-core/src/test/java/datadog/trace/core/otlp/common/OtlpResourceJsonTest.java index e73311993cf..102027cf668 100644 --- a/dd-trace-core/src/test/java/datadog/trace/core/otlp/common/OtlpResourceJsonTest.java +++ b/dd-trace-core/src/test/java/datadog/trace/core/otlp/common/OtlpResourceJsonTest.java @@ -5,6 +5,7 @@ import static datadog.trace.api.config.GeneralConfig.EXPERIMENTAL_PROPAGATE_PROCESS_TAGS_ENABLED; import static datadog.trace.api.config.GeneralConfig.SERVICE_NAME; import static datadog.trace.api.config.GeneralConfig.TAGS; +import static datadog.trace.api.config.GeneralConfig.TRACE_OTEL_SEMANTICS_ENABLED; import static datadog.trace.api.config.GeneralConfig.VERSION; import static datadog.trace.api.config.OtlpConfig.OTEL_TRACES_SPAN_METRICS_ENABLED; import static datadog.trace.api.config.TracerConfig.TRACE_REPORT_HOSTNAME; @@ -125,7 +126,9 @@ static Stream resourceFragmentCases() { + "SERVICE.VERSION:ignored-version," + "telemetry.sdk.name:ignored-sdk," + "telemetry.sdk.version:ignored-version," - + "telemetry.sdk.language:ignored-language"), + + "telemetry.sdk.language:ignored-language," + + "datadog.sdk.semantics:ignored-semantics," + + "_dd.sdk.otlp_export:ignored-export"), attrs( "service.name", "my-service", "deployment.environment.name", "staging", @@ -147,6 +150,32 @@ void testBuildResourceFragment( assertEquals(expectedAttributes, actualAttributes, "For case: " + caseName); } + /** + * On the OTLP path the export-mode marker is a resource attribute and is always {@code "true"} -- + * reaching this encoder means the payload is leaving over OTLP. + */ + @Test + void traceResourceAttributesCarryOtlpExportMarker() throws IOException { + Config config = Config.get(props(SERVICE_NAME, "my-service")); + + Map attributes = + parseResourceAttributes( + OtlpResourceJson.buildResourceFragment(config, traceResourceAttributes(config))); + + assertEquals("true", attributes.get("_dd.sdk.otlp_export")); + } + + @Test + void usesOtelSdkSemanticsWhenEnabled() throws IOException { + Config config = Config.get(props(TRACE_OTEL_SEMANTICS_ENABLED, "true")); + + Map attributes = + parseResourceAttributes( + OtlpResourceJson.buildResourceFragment(config, traceResourceAttributes(config))); + + assertEquals("otel", attributes.get("datadog.sdk.semantics")); + } + /** The datadog-attrs variant carries {@code datadog.runtime_id}; the plain variant omits it. */ @Test void datadogResourceAttributesVariantCarriesRuntimeId() throws IOException { @@ -211,6 +240,8 @@ void statsComputedVariantCarriesMarker() throws IOException { assertEquals( "true", withMarker.get("_dd.stats_computed"), "marker present when stats computed"); assertFalse(without.containsKey("_dd.stats_computed"), "marker absent when stats not computed"); + assertEquals("datadog", without.get("datadog.sdk.semantics")); + assertEquals("true", without.get("_dd.sdk.otlp_export")); } @Test diff --git a/dd-trace-core/src/test/java/datadog/trace/core/otlp/common/OtlpResourceProtoTest.java b/dd-trace-core/src/test/java/datadog/trace/core/otlp/common/OtlpResourceProtoTest.java index 24dda2ac60a..f2831b3bd19 100644 --- a/dd-trace-core/src/test/java/datadog/trace/core/otlp/common/OtlpResourceProtoTest.java +++ b/dd-trace-core/src/test/java/datadog/trace/core/otlp/common/OtlpResourceProtoTest.java @@ -5,6 +5,7 @@ import static datadog.trace.api.config.GeneralConfig.EXPERIMENTAL_PROPAGATE_PROCESS_TAGS_ENABLED; import static datadog.trace.api.config.GeneralConfig.SERVICE_NAME; import static datadog.trace.api.config.GeneralConfig.TAGS; +import static datadog.trace.api.config.GeneralConfig.TRACE_OTEL_SEMANTICS_ENABLED; import static datadog.trace.api.config.GeneralConfig.VERSION; import static datadog.trace.api.config.OtlpConfig.OTEL_TRACES_SPAN_METRICS_ENABLED; import static datadog.trace.api.config.TracerConfig.TRACE_REPORT_HOSTNAME; @@ -147,7 +148,9 @@ static Stream resourceMessageCases() { + "SERVICE.VERSION:ignored-version," + "telemetry.sdk.name:ignored-sdk," + "telemetry.sdk.version:ignored-version," - + "telemetry.sdk.language:ignored-language"), + + "telemetry.sdk.language:ignored-language," + + "datadog.sdk.semantics:ignored-semantics," + + "_dd.sdk.otlp_export:ignored-export"), attrs( "service.name", "my-service", "deployment.environment.name", "staging", @@ -169,6 +172,32 @@ void testBuildResourceMessage( assertEquals(expectedAttributes, actualAttributes, "For case: " + caseName); } + /** + * On the OTLP path the export-mode marker is a resource attribute and is always {@code "true"} -- + * reaching this encoder means the payload is leaving over OTLP. + */ + @Test + void traceResourceAttributesCarryOtlpExportMarker() throws IOException { + Config config = Config.get(props(SERVICE_NAME, "my-service")); + + Map attributes = + parseResourceAttributes( + OtlpResourceProto.buildResourceMessage(config, traceResourceAttributes(config))); + + assertEquals("true", attributes.get("_dd.sdk.otlp_export")); + } + + @Test + void usesOtelSdkSemanticsWhenEnabled() throws IOException { + Config config = Config.get(props(TRACE_OTEL_SEMANTICS_ENABLED, "true")); + + Map attributes = + parseResourceAttributes( + OtlpResourceProto.buildResourceMessage(config, traceResourceAttributes(config))); + + assertEquals("otel", attributes.get("datadog.sdk.semantics")); + } + /** * The datadog-attrs variant ({@code buildResourceMessage(config, datadogResourceAttributes)}) * carries {@code datadog.runtime_id}; the plain variant omits it. @@ -236,6 +265,8 @@ void statsComputedVariantCarriesMarker() throws IOException { assertEquals( "true", withMarker.get("_dd.stats_computed"), "marker present when stats computed"); assertFalse(without.containsKey("_dd.stats_computed"), "marker absent when stats not computed"); + assertEquals("datadog", without.get("datadog.sdk.semantics")); + assertEquals("true", without.get("_dd.sdk.otlp_export")); } // ── parsing helpers ─────────────────────────────────────────────────────── diff --git a/dd-trace-core/src/traceAgentTest/java/TraceGenerator.java b/dd-trace-core/src/traceAgentTest/java/TraceGenerator.java index 1e1a9c58fe3..da0911022e9 100644 --- a/dd-trace-core/src/traceAgentTest/java/TraceGenerator.java +++ b/dd-trace-core/src/traceAgentTest/java/TraceGenerator.java @@ -6,6 +6,7 @@ import static java.util.Collections.emptyList; import static java.util.Collections.emptyMap; +import datadog.trace.api.Config; import datadog.trace.api.DDSpanId; import datadog.trace.api.DDTags; import datadog.trace.api.DDTraceId; @@ -182,6 +183,7 @@ static class PojoSpan implements CoreSpan { null, 0, getTagsForSerialization(), + UTF8BytesString.create(String.valueOf(Config.get().isOtlpTracesExportEnabled())), emptyList()); }