From b47d44a365390fe96f234e800d9cbea550d266c4 Mon Sep 17 00:00:00 2001 From: Sarah Chen Date: Tue, 28 Jul 2026 11:30:15 -0400 Subject: [PATCH 1/2] Propagate sampling metadata across trace chunks --- .../java/datadog/trace/core/CoreTracer.java | 3 + .../datadog/trace/core/DDSpanContext.java | 4 + ...panSamplingMechanismSerializationTest.java | 156 ++++++++++++++++++ 3 files changed, 163 insertions(+) create mode 100644 dd-trace-core/src/test/java/datadog/trace/core/LateSpanSamplingMechanismSerializationTest.java diff --git a/dd-trace-core/src/main/java/datadog/trace/core/CoreTracer.java b/dd-trace-core/src/main/java/datadog/trace/core/CoreTracer.java index 289647c5c74..178042f8669 100644 --- a/dd-trace-core/src/main/java/datadog/trace/core/CoreTracer.java +++ b/dd-trace-core/src/main/java/datadog/trace/core/CoreTracer.java @@ -1310,6 +1310,9 @@ void write(final SpanList trace) { spanToSample.forceKeep(forceKeep); boolean published = forceKeep || traceCollector.sample(spanToSample); if (published) { + if (rootSpan != null && writtenTrace.get(0) != rootSpan) { + writtenTrace.get(0).spanContext().copyDecisionMakerFrom(rootSpan.spanContext()); + } if (!apmTracingEnabled) { // Stamp the billing marker on every span of each exported chunk so the intake does not bill // APM host usage. 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..e30647699e6 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 @@ -1393,6 +1393,10 @@ public PropagationTags getPropagationTags() { return getRootSpanContextOrThis().propagationTags; } + void copyDecisionMakerFrom(DDSpanContext source) { + propagationTags.updateAndLockDecisionMaker(source.getPropagationTags()); + } + /** TraceSegment Implementation */ @Override public void setTagTop(String key, Object value, boolean sanitize) { diff --git a/dd-trace-core/src/test/java/datadog/trace/core/LateSpanSamplingMechanismSerializationTest.java b/dd-trace-core/src/test/java/datadog/trace/core/LateSpanSamplingMechanismSerializationTest.java new file mode 100644 index 00000000000..aab4b5f2120 --- /dev/null +++ b/dd-trace-core/src/test/java/datadog/trace/core/LateSpanSamplingMechanismSerializationTest.java @@ -0,0 +1,156 @@ +package datadog.trace.core; + +import static datadog.trace.api.sampling.PrioritySampling.USER_KEEP; +import static datadog.trace.api.sampling.SamplingMechanism.REMOTE_ADAPTIVE_RULE; +import static datadog.trace.common.sampling.RuleBasedTraceSampler.SAMPLING_RULE_RATE; +import static java.util.Collections.singletonList; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import datadog.communication.serialization.ByteBufferConsumer; +import datadog.communication.serialization.FlushingBuffer; +import datadog.communication.serialization.msgpack.MsgPackWriter; +import datadog.trace.common.writer.ListWriter; +import datadog.trace.common.writer.ddagent.TraceMapperV0_4; +import java.nio.ByteBuffer; +import java.util.ArrayList; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.concurrent.TimeoutException; +import org.junit.jupiter.api.Test; +import org.msgpack.core.MessagePack; +import org.msgpack.core.MessageUnpacker; + +/** + * Verifies that trace-level sampling metadata is preserved when a root and a later child are + * exported in separate chunks (e.g. during partial sampling). + */ +class LateSpanSamplingMechanismSerializationTest extends DDCoreJavaSpecification { + + @Override + protected boolean useStrictTraceWrites() { + return false; + } + + @Test + void rootlessLateChunkRetainsAdaptiveSamplingMechanism() + throws InterruptedException, TimeoutException { + SerializingWriter writer = new SerializingWriter(); + CoreTracer tracer = + tracerBuilder() + .writer(writer) + // Keep the partial-flush threshold above the trace size so the pending-trace idle + // timeout controls the first write. + .partialFlushMinSpans(1000) + .build(); + + try { + DDSpan root = (DDSpan) tracer.buildSpan("test", "root").start(); + DDSpan lateChild = + (DDSpan) tracer.buildSpan("test", "late-child").asChildOf(root.spanContext()).start(); + PendingTrace trace = (PendingTrace) root.spanContext().getTraceCollector(); + + root.setSamplingPriority(USER_KEEP, REMOTE_ADAPTIVE_RULE); + root.setMetric(SAMPLING_RULE_RATE, 0.25); + + // Root finishes while the child is still running -> root is buffered, not yet written. + root.finish(); + assertTrue(writer.isEmpty()); + + // Write the buffered root as the PendingTraceBuffer worker does after its idle timeout. + trace.write(); + writer.waitForTraces(1); + + // The child finishes after the root was written, so it is exported in a root-less chunk. + lateChild.finish(); + trace.write(); + writer.waitForTraces(2); + + assertEquals(singletonList(root), writer.get(0)); + assertEquals(singletonList(lateChild), writer.get(1)); + + SerializedChunk rootChunk = writer.serializedChunks.get(0); + SerializedChunk lateChunk = writer.serializedChunks.get(1); + + assertEquals( + (int) USER_KEEP, rootChunk.metrics.get(DDSpanContext.PRIORITY_SAMPLING_KEY).intValue()); + assertEquals("-12", rootChunk.meta.get("_dd.p.dm")); + + assertEquals( + (int) USER_KEEP, lateChunk.metrics.get(DDSpanContext.PRIORITY_SAMPLING_KEY).intValue()); + // The late chunk's decision maker must match its propagated trace-level sampling priority. + assertEquals("-12", lateChunk.meta.get("_dd.p.dm")); + } finally { + tracer.close(); + } + } + + private static final class SerializingWriter extends ListWriter { + private final List serializedChunks = new ArrayList<>(); + + @Override + public void write(List trace) { + serializedChunks.add(serialize(trace)); + super.write(trace); + } + + private static SerializedChunk serialize(List trace) { + TraceMapperV0_4 mapper = new TraceMapperV0_4(); + CaptureBuffer capture = new CaptureBuffer(); + MsgPackWriter packer = new MsgPackWriter(new FlushingBuffer(16 * 1024, capture)); + assertTrue(packer.format(trace, mapper)); + packer.flush(); + + try (MessageUnpacker unpacker = MessagePack.newDefaultUnpacker(capture.bytes)) { + assertEquals(1, unpacker.unpackArrayHeader()); + int fieldCount = unpacker.unpackMapHeader(); + Map metrics = new HashMap<>(); + Map meta = new HashMap<>(); + + for (int i = 0; i < fieldCount; i++) { + String field = unpacker.unpackString(); + if ("metrics".equals(field)) { + int size = unpacker.unpackMapHeader(); + for (int j = 0; j < size; j++) { + metrics.put(unpacker.unpackString(), unpacker.unpackValue().asNumberValue().toInt()); + } + } else if ("meta".equals(field)) { + int size = unpacker.unpackMapHeader(); + for (int j = 0; j < size; j++) { + meta.put(unpacker.unpackString(), unpacker.unpackString()); + } + } else { + unpacker.unpackValue(); + } + } + return new SerializedChunk(metrics, meta); + } catch (Exception e) { + throw new IllegalStateException("Unable to decode serialized trace chunk", e); + } finally { + mapper.reset(); + } + } + } + + private static final class CaptureBuffer implements ByteBufferConsumer { + private byte[] bytes; + + @Override + public void accept(int messageCount, ByteBuffer buffer) { + assertEquals(1, messageCount); + bytes = new byte[buffer.remaining()]; + buffer.get(bytes); + } + } + + private static final class SerializedChunk { + private final Map metrics; + private final Map meta; + + private SerializedChunk(Map metrics, Map meta) { + this.metrics = metrics; + this.meta = meta; + } + } +} From 71f3bbb770bf488af3e9029c1fb5de8cc77320ba Mon Sep 17 00:00:00 2001 From: Sarah Chen Date: Tue, 28 Jul 2026 11:49:31 -0400 Subject: [PATCH 2/2] Simplify test --- ...panSamplingMechanismSerializationTest.java | 100 +++++++----------- 1 file changed, 37 insertions(+), 63 deletions(-) diff --git a/dd-trace-core/src/test/java/datadog/trace/core/LateSpanSamplingMechanismSerializationTest.java b/dd-trace-core/src/test/java/datadog/trace/core/LateSpanSamplingMechanismSerializationTest.java index aab4b5f2120..f68ed798494 100644 --- a/dd-trace-core/src/test/java/datadog/trace/core/LateSpanSamplingMechanismSerializationTest.java +++ b/dd-trace-core/src/test/java/datadog/trace/core/LateSpanSamplingMechanismSerializationTest.java @@ -7,20 +7,20 @@ import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertTrue; -import datadog.communication.serialization.ByteBufferConsumer; -import datadog.communication.serialization.FlushingBuffer; +import datadog.communication.serialization.GrowableBuffer; import datadog.communication.serialization.msgpack.MsgPackWriter; import datadog.trace.common.writer.ListWriter; import datadog.trace.common.writer.ddagent.TraceMapperV0_4; -import java.nio.ByteBuffer; +import java.io.IOException; import java.util.ArrayList; -import java.util.HashMap; import java.util.List; -import java.util.Map; import java.util.concurrent.TimeoutException; import org.junit.jupiter.api.Test; import org.msgpack.core.MessagePack; import org.msgpack.core.MessageUnpacker; +import org.msgpack.value.MapValue; +import org.msgpack.value.Value; +import org.msgpack.value.ValueFactory; /** * Verifies that trace-level sampling metadata is preserved when a root and a later child are @@ -70,87 +70,61 @@ void rootlessLateChunkRetainsAdaptiveSamplingMechanism() assertEquals(singletonList(root), writer.get(0)); assertEquals(singletonList(lateChild), writer.get(1)); - SerializedChunk rootChunk = writer.serializedChunks.get(0); - SerializedChunk lateChunk = writer.serializedChunks.get(1); + MapValue rootChunk = writer.serializedSpans.get(0); + MapValue lateChunk = writer.serializedSpans.get(1); assertEquals( - (int) USER_KEEP, rootChunk.metrics.get(DDSpanContext.PRIORITY_SAMPLING_KEY).intValue()); - assertEquals("-12", rootChunk.meta.get("_dd.p.dm")); + (int) USER_KEEP, + numberValue(rootChunk, "metrics", DDSpanContext.PRIORITY_SAMPLING_KEY).intValue()); + assertEquals("-12", stringValue(rootChunk, "meta", "_dd.p.dm")); assertEquals( - (int) USER_KEEP, lateChunk.metrics.get(DDSpanContext.PRIORITY_SAMPLING_KEY).intValue()); + (int) USER_KEEP, + numberValue(lateChunk, "metrics", DDSpanContext.PRIORITY_SAMPLING_KEY).intValue()); // The late chunk's decision maker must match its propagated trace-level sampling priority. - assertEquals("-12", lateChunk.meta.get("_dd.p.dm")); + assertEquals("-12", stringValue(lateChunk, "meta", "_dd.p.dm")); } finally { tracer.close(); } } + private static Number numberValue(MapValue span, String mapName, String key) { + Value value = value(span, mapName, key); + return value == null ? null : value.asNumberValue().toInt(); + } + + private static String stringValue(MapValue span, String mapName, String key) { + Value value = value(span, mapName, key); + return value == null ? null : value.asStringValue().asString(); + } + + private static Value value(MapValue span, String mapName, String key) { + Value map = span.map().get(ValueFactory.newString(mapName)); + return map.asMapValue().map().get(ValueFactory.newString(key)); + } + private static final class SerializingWriter extends ListWriter { - private final List serializedChunks = new ArrayList<>(); + private final List serializedSpans = new ArrayList<>(); @Override public void write(List trace) { - serializedChunks.add(serialize(trace)); + serializedSpans.add(serialize(trace)); super.write(trace); } - private static SerializedChunk serialize(List trace) { + private static MapValue serialize(List trace) { TraceMapperV0_4 mapper = new TraceMapperV0_4(); - CaptureBuffer capture = new CaptureBuffer(); - MsgPackWriter packer = new MsgPackWriter(new FlushingBuffer(16 * 1024, capture)); + GrowableBuffer buffer = new GrowableBuffer(16 * 1024); + MsgPackWriter packer = new MsgPackWriter(buffer); assertTrue(packer.format(trace, mapper)); - packer.flush(); - - try (MessageUnpacker unpacker = MessagePack.newDefaultUnpacker(capture.bytes)) { - assertEquals(1, unpacker.unpackArrayHeader()); - int fieldCount = unpacker.unpackMapHeader(); - Map metrics = new HashMap<>(); - Map meta = new HashMap<>(); - - for (int i = 0; i < fieldCount; i++) { - String field = unpacker.unpackString(); - if ("metrics".equals(field)) { - int size = unpacker.unpackMapHeader(); - for (int j = 0; j < size; j++) { - metrics.put(unpacker.unpackString(), unpacker.unpackValue().asNumberValue().toInt()); - } - } else if ("meta".equals(field)) { - int size = unpacker.unpackMapHeader(); - for (int j = 0; j < size; j++) { - meta.put(unpacker.unpackString(), unpacker.unpackString()); - } - } else { - unpacker.unpackValue(); - } - } - return new SerializedChunk(metrics, meta); - } catch (Exception e) { + + try (MessageUnpacker unpacker = MessagePack.newDefaultUnpacker(buffer.slice())) { + return unpacker.unpackValue().asArrayValue().get(0).asMapValue(); + } catch (IOException e) { throw new IllegalStateException("Unable to decode serialized trace chunk", e); } finally { mapper.reset(); } } } - - private static final class CaptureBuffer implements ByteBufferConsumer { - private byte[] bytes; - - @Override - public void accept(int messageCount, ByteBuffer buffer) { - assertEquals(1, messageCount); - bytes = new byte[buffer.remaining()]; - buffer.get(bytes); - } - } - - private static final class SerializedChunk { - private final Map metrics; - private final Map meta; - - private SerializedChunk(Map metrics, Map meta) { - this.metrics = metrics; - this.meta = meta; - } - } }