diff --git a/dd-java-agent/instrumentation/opentelemetry/opentelemetry-0.3/src/test/groovy/OpenTelemetryTest.groovy b/dd-java-agent/instrumentation/opentelemetry/opentelemetry-0.3/src/test/groovy/OpenTelemetryTest.groovy index 62f315e6191..e222e462a83 100644 --- a/dd-java-agent/instrumentation/opentelemetry/opentelemetry-0.3/src/test/groovy/OpenTelemetryTest.groovy +++ b/dd-java-agent/instrumentation/opentelemetry/opentelemetry-0.3/src/test/groovy/OpenTelemetryTest.groovy @@ -287,6 +287,7 @@ class OpenTelemetryTest extends InstrumentationSpecification { } if (contextPriority == UNSET) { expectedTracestate += ";t.ksr:1" + expectedTracestate += ",ot=${span.delegate.spanContext().propagationTags.samplingState().otelTraceState}" } if (traceId.toHighOrderLong() != 0) { expectedDataTags << "_dd.p.tid=" + traceId.toHexStringPadded(32).substring(0, 16) diff --git a/dd-java-agent/instrumentation/opentracing/opentracing-0.31/src/test/groovy/OpenTracing31Test.groovy b/dd-java-agent/instrumentation/opentracing/opentracing-0.31/src/test/groovy/OpenTracing31Test.groovy index 53f4d0a62fc..a9bcbd17f0c 100644 --- a/dd-java-agent/instrumentation/opentracing/opentracing-0.31/src/test/groovy/OpenTracing31Test.groovy +++ b/dd-java-agent/instrumentation/opentracing/opentracing-0.31/src/test/groovy/OpenTracing31Test.groovy @@ -295,6 +295,7 @@ class OpenTracing31Test extends InstrumentationSpecification { } if (contextPriority == UNSET) { expectedTracestate += ";t.ksr:1" + expectedTracestate += ",ot=${context.delegate.propagationTags.samplingState().otelTraceState}" datadogTags << "_dd.p.ksr=1" } def expectedTextMap = [ diff --git a/dd-java-agent/instrumentation/opentracing/opentracing-0.32/src/test/groovy/OpenTracing32Test.groovy b/dd-java-agent/instrumentation/opentracing/opentracing-0.32/src/test/groovy/OpenTracing32Test.groovy index fe452134938..0709d70f832 100644 --- a/dd-java-agent/instrumentation/opentracing/opentracing-0.32/src/test/groovy/OpenTracing32Test.groovy +++ b/dd-java-agent/instrumentation/opentracing/opentracing-0.32/src/test/groovy/OpenTracing32Test.groovy @@ -311,6 +311,7 @@ class OpenTracing32Test extends InstrumentationSpecification { } if (contextPriority == UNSET) { expectedTracestate+= ";t.ksr:1" + expectedTracestate+= ",ot=${context.delegate.propagationTags.samplingState().otelTraceState}" datadogTags << "_dd.p.ksr=1" } def expectedTextMap = [ diff --git a/dd-trace-core/src/main/java/datadog/trace/common/sampling/RateByServiceTraceSampler.java b/dd-trace-core/src/main/java/datadog/trace/common/sampling/RateByServiceTraceSampler.java index 0940469c61c..ab24b2200b0 100644 --- a/dd-trace-core/src/main/java/datadog/trace/common/sampling/RateByServiceTraceSampler.java +++ b/dd-trace-core/src/main/java/datadog/trace/common/sampling/RateByServiceTraceSampler.java @@ -59,19 +59,14 @@ public > void setSamplingPriority(final T span) { final RateSamplersByEnvAndService rates = serviceRates; RateSampler sampler = rates.getSampler(env, serviceName); - if (sampler.sample(span)) { - span.setSamplingPriority( - PrioritySampling.SAMPLER_KEEP, - SAMPLING_AGENT_RATE, - sampler.getSampleRate(), - SamplingMechanism.AGENT_RATE); - } else { - span.setSamplingPriority( - PrioritySampling.SAMPLER_DROP, - SAMPLING_AGENT_RATE, - sampler.getSampleRate(), - SamplingMechanism.AGENT_RATE); - } + boolean sampled = sampler.sample(span); + int samplingPriority = sampled ? PrioritySampling.SAMPLER_KEEP : PrioritySampling.SAMPLER_DROP; + span.setSamplingPriority( + samplingPriority, + SAMPLING_AGENT_RATE, + sampler.getSampleRate(), + sampled, + SamplingMechanism.AGENT_RATE); } private > String getSpanEnv(final T span) { diff --git a/dd-trace-core/src/main/java/datadog/trace/common/sampling/RuleBasedTraceSampler.java b/dd-trace-core/src/main/java/datadog/trace/common/sampling/RuleBasedTraceSampler.java index 1746a1a3c54..c66159607b1 100644 --- a/dd-trace-core/src/main/java/datadog/trace/common/sampling/RuleBasedTraceSampler.java +++ b/dd-trace-core/src/main/java/datadog/trace/common/sampling/RuleBasedTraceSampler.java @@ -146,18 +146,21 @@ public > void setSamplingPriority(final T span) { if (matchedRule == null) { fallbackSampler.setSamplingPriority(span); } else { - if (matchedRule.sample(span)) { + boolean sampled = matchedRule.sample(span); + if (sampled) { if (rateLimiter.tryAcquire()) { span.setSamplingPriority( PrioritySampling.USER_KEEP, SAMPLING_RULE_RATE, matchedRule.getSampler().getSampleRate(), + true, matchedRule.getMechanism()); } else { span.setSamplingPriority( PrioritySampling.USER_DROP, SAMPLING_RULE_RATE, matchedRule.getSampler().getSampleRate(), + true, matchedRule.getMechanism()); } span.setMetric(SAMPLING_LIMIT_RATE, rateLimit); @@ -166,6 +169,7 @@ public > void setSamplingPriority(final T span) { PrioritySampling.USER_DROP, SAMPLING_RULE_RATE, matchedRule.getSampler().getSampleRate(), + false, matchedRule.getMechanism()); } } 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 b2ab55c8e25..f7b5bcb7c3a 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 @@ -122,7 +122,11 @@ default void processTagsAndBaggageWithStructuredLinks( T setSamplingPriority(int samplingPriority, int samplingMechanism); T setSamplingPriority( - int samplingPriority, CharSequence rate, double sampleRate, int samplingMechanism); + int samplingPriority, + CharSequence rate, + double sampleRate, + boolean probabilitySamplingResult, + int samplingMechanism); T setSpanSamplingPriority(double rate, int limit); 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 57809e76069..fa94775f214 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 @@ -647,15 +647,18 @@ public final DDSpan setSamplingPriority(final int newPriority, int samplingMecha @Override public DDSpan setSamplingPriority( - int samplingPriority, CharSequence rate, double sampleRate, int samplingMechanism) { - if (context.setSamplingPriority(samplingPriority, samplingMechanism)) { + int samplingPriority, + CharSequence rate, + double sampleRate, + boolean probabilitySamplingResult, + int samplingMechanism) { + if (context.setSamplingPriority( + samplingPriority, + samplingMechanism, + sampleRate, + probabilitySamplingResult, + getTraceId().toLong())) { setMetric(rate, sampleRate); - if (samplingMechanism == SamplingMechanism.AGENT_RATE - || samplingMechanism == SamplingMechanism.LOCAL_USER_RULE - || samplingMechanism == SamplingMechanism.REMOTE_USER_RULE - || samplingMechanism == SamplingMechanism.REMOTE_ADAPTIVE_RULE) { - context.getPropagationTags().updateKnuthSamplingRate(sampleRate); - } } return this; } 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..6359e8ece2d 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 @@ -46,7 +46,6 @@ import java.util.TreeMap; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ThreadLocalRandom; -import java.util.concurrent.atomic.AtomicIntegerFieldUpdater; import java.util.function.Function; import javax.annotation.Nonnull; import org.slf4j.Logger; @@ -161,11 +160,6 @@ public class DDSpanContext private volatile boolean topLevel; - private static final AtomicIntegerFieldUpdater SAMPLING_PRIORITY_UPDATER = - AtomicIntegerFieldUpdater.newUpdater(DDSpanContext.class, "samplingPriority"); - - private volatile int samplingPriority = PrioritySampling.UNSET; - /** The origin of the trace. (eg. Synthetics, CI App) */ private volatile CharSequence origin; @@ -667,10 +661,6 @@ public void forceKeep(byte samplingMechanism) { } private void forceKeepThisSpan(byte samplingMechanism) { - // if the user really wants to keep this trace chunk, we will let them, - // even if the old sampling priority and mechanism have already propagated - SAMPLING_PRIORITY_UPDATER.set(this, PrioritySampling.USER_KEEP); - // record force keep decision for future distributed trace propagation propagationTags.forceKeep(samplingMechanism); } @@ -710,26 +700,43 @@ private boolean setThisSpanSamplingPriority(final int newPriority, final int new if (!validateSamplingPriority(newPriority, newMechanism)) { return false; } - if (SamplingMechanism.canAvoidSamplingPriorityLock(newPriority, newMechanism)) { - SAMPLING_PRIORITY_UPDATER.set(this, newPriority); - propagationTags.updateTraceSamplingPriority(newPriority, newMechanism); - return true; - } - if (!SAMPLING_PRIORITY_UPDATER.compareAndSet(this, PrioritySampling.UNSET, newPriority)) { + boolean updated = + propagationTags.tryUpdateTraceSamplingPriority( + newPriority, + newMechanism, + SamplingMechanism.canAvoidSamplingPriorityLock(newPriority, newMechanism)); + if (!updated) { if (log.isDebugEnabled()) { log.debug( "samplingPriority locked at priority: {}. Refusing to set to priority: {} mechanism: {}", - samplingPriority, + propagationTags.getSamplingPriority(), newPriority, newMechanism); } return false; } - // set trace level sampling priority tag propagationTags - propagationTags.updateTraceSamplingPriority(newPriority, newMechanism); return true; } + public boolean setSamplingPriority( + final int newPriority, + final int newMechanism, + final double sampleRate, + final boolean probabilitySamplingResult, + final long traceIdLowOrderBits) { + DDSpanContext spanContext = getRootSpanContextOrThis(); + if (!spanContext.validateSamplingPriority(newPriority, newMechanism)) { + return false; + } + return spanContext.propagationTags.tryUpdateProbabilitySamplingDecision( + newPriority, + newMechanism, + sampleRate, + probabilitySamplingResult, + traceIdLowOrderBits, + SamplingMechanism.canAvoidSamplingPriorityLock(newPriority, newMechanism)); + } + private boolean validateSamplingPriority(final int newPriority, final int newMechanism) { if (newPriority == PrioritySampling.UNSET) { log.debug("{}: Refusing to set samplingPriority to UNSET", this); @@ -758,7 +765,7 @@ private boolean validateSamplingPriority(final int newPriority, final int newMec @Override public int getSamplingPriority() { - return getRootSpanContextOrThis().samplingPriority; + return getRootSpanContextOrThis().propagationTags.getSamplingPriority(); } public void setSpanSamplingPriority(double rate, int limit) { @@ -790,7 +797,7 @@ public boolean lockSamplingPriority() { return rootSpan.spanContext().lockSamplingPriority(); } - return SAMPLING_PRIORITY_UPDATER.get(this) != PrioritySampling.UNSET; + return propagationTags.getSamplingPriority() != PrioritySampling.UNSET; } public CharSequence getOrigin() { @@ -1258,8 +1265,11 @@ public TagMap getTags() { tags.put(DDTags.THREAD_ID, threadId); // maintain previously observable type of the thread name :| tags.put(DDTags.THREAD_NAME, threadName.toString()); - if (samplingPriority != PrioritySampling.UNSET) { - tags.put(SAMPLE_RATE_KEY, samplingPriority); + int currentSamplingPriority = getSamplingPriority(); + // add _sample_rate tag only on the root/local span owning the decision + if (getRootSpanContextIfDifferent() == null + && currentSamplingPriority != PrioritySampling.UNSET) { + tags.put(SAMPLE_RATE_KEY, currentSamplingPriority); } if (httpStatusCode != 0) { tags.put(Tags.HTTP_STATUS, (int) httpStatusCode); @@ -1408,7 +1418,7 @@ void processTagsAndBaggage( threadName, unsafeTags, baggageItemsWithPropagationTags, - samplingPriority != PrioritySampling.UNSET ? samplingPriority : getSamplingPriority(), + getSamplingPriority(), measured, topLevel, httpStatusCode, diff --git a/dd-trace-core/src/main/java/datadog/trace/core/otlp/trace/OtlpTraceJson.java b/dd-trace-core/src/main/java/datadog/trace/core/otlp/trace/OtlpTraceJson.java index a56498d0a23..845f870b965 100644 --- a/dd-trace-core/src/main/java/datadog/trace/core/otlp/trace/OtlpTraceJson.java +++ b/dd-trace-core/src/main/java/datadog/trace/core/otlp/trace/OtlpTraceJson.java @@ -34,6 +34,7 @@ import datadog.trace.core.MetadataConsumer; import datadog.trace.core.PendingTrace; import datadog.trace.core.propagation.PropagationTags; +import datadog.trace.core.propagation.PropagationTags.SamplingState; import java.util.List; import java.util.Map; @@ -51,13 +52,14 @@ private OtlpTraceJson() {} public static void writeSpan( JsonWriter writer, DDSpan span, MetaWriter metaWriter, List links) { PropagationTags propagationTags = span.spanContext().getPropagationTags(); + SamplingState samplingState = propagationTags.samplingState(); writer.beginObject(); writer.name("traceId").value(hexTraceId(span.getTraceId())); writer.name("spanId").value(hexSpanId(span.getSpanId())); - String tracestate = propagationTags.getW3CTracestate(); + String tracestate = propagationTags.getW3CTracestate(samplingState); if (tracestate != null) { writer.name("traceState").value(tracestate); } @@ -67,7 +69,7 @@ public static void writeSpan( } int traceFlags = NO_TRACE_FLAGS; - if (span.samplingPriority() > 0) { + if (samplingState.getSamplingPriority() > 0) { traceFlags |= SAMPLED_TRACE_FLAG; } if (span.spanContext().isRemote()) { diff --git a/dd-trace-core/src/main/java/datadog/trace/core/otlp/trace/OtlpTraceProto.java b/dd-trace-core/src/main/java/datadog/trace/core/otlp/trace/OtlpTraceProto.java index 258cfb73669..707225cf5a6 100644 --- a/dd-trace-core/src/main/java/datadog/trace/core/otlp/trace/OtlpTraceProto.java +++ b/dd-trace-core/src/main/java/datadog/trace/core/otlp/trace/OtlpTraceProto.java @@ -48,6 +48,7 @@ import datadog.trace.core.PendingTrace; import datadog.trace.core.otlp.common.OtlpProtoBuffer; import datadog.trace.core.propagation.PropagationTags; +import datadog.trace.core.propagation.PropagationTags.SamplingState; /** Provides optimized writers for OpenTelemetry's "trace.proto" wire protocol. */ public final class OtlpTraceProto { @@ -84,6 +85,7 @@ public static int recordSpanMessage( int nestedSpanLinkBytes, OtlpProtoBuffer protobuf) { PropagationTags propagationTags = span.spanContext().getPropagationTags(); + SamplingState samplingState = propagationTags.samplingState(); writeTag(buf, 1, LEN_WIRE_TYPE); writeTraceId(buf, span.getTraceId()); @@ -91,7 +93,7 @@ public static int recordSpanMessage( writeTag(buf, 2, LEN_WIRE_TYPE); writeSpanId(buf, span.getSpanId()); - String tracestate = propagationTags.getW3CTracestate(); + String tracestate = propagationTags.getW3CTracestate(samplingState); if (tracestate != null) { writeTag(buf, 3, LEN_WIRE_TYPE); writeString(buf, tracestate); @@ -103,7 +105,7 @@ public static int recordSpanMessage( } int traceFlags = NO_TRACE_FLAGS; - if (span.samplingPriority() > 0) { + if (samplingState.getSamplingPriority() > 0) { traceFlags |= SAMPLED_TRACE_FLAG; } if (span.spanContext().isRemote()) { diff --git a/dd-trace-core/src/main/java/datadog/trace/core/propagation/PropagationTags.java b/dd-trace-core/src/main/java/datadog/trace/core/propagation/PropagationTags.java index 0ebe630c87a..e6f406c4f0e 100644 --- a/dd-trace-core/src/main/java/datadog/trace/core/propagation/PropagationTags.java +++ b/dd-trace-core/src/main/java/datadog/trace/core/propagation/PropagationTags.java @@ -22,6 +22,47 @@ */ public abstract class PropagationTags { + public static final class SamplingState { + private final int samplingPriority; + private final String tracestate; + private final CharSequence otelTraceState; + private final CharSequence decisionMaker; + private final CharSequence knuthSamplingRate; + + public SamplingState( + int samplingPriority, + String tracestate, + CharSequence otelTraceState, + CharSequence decisionMaker, + CharSequence knuthSamplingRate) { + this.samplingPriority = samplingPriority; + this.tracestate = tracestate; + this.otelTraceState = otelTraceState; + this.decisionMaker = decisionMaker; + this.knuthSamplingRate = knuthSamplingRate; + } + + public int getSamplingPriority() { + return samplingPriority; + } + + public String getTracestate() { + return tracestate; + } + + public CharSequence getOtelTraceState() { + return otelTraceState; + } + + public CharSequence getDecisionMaker() { + return decisionMaker; + } + + public CharSequence getKnuthSamplingRate() { + return knuthSamplingRate; + } + } + public static PropagationTags.Factory factory(Config config) { return factory(config.getxDatadogTagsMaxLength()); } @@ -65,10 +106,23 @@ public interface Factory { */ public abstract void updateTraceSamplingPriority(int samplingPriority, int samplingMechanism); + public abstract boolean tryUpdateTraceSamplingPriority( + int samplingPriority, int samplingMechanism, boolean allowOverride); + + public abstract boolean tryUpdateProbabilitySamplingDecision( + int samplingPriority, + int samplingMechanism, + double sampleRate, + boolean probabilitySamplingResult, + long traceIdLowOrderBits, + boolean allowOverride); + public abstract void forceKeep(int samplingMechanism); public abstract int getSamplingPriority(); + public abstract SamplingState samplingState(); + public abstract void updateTraceOrigin(CharSequence origin); public abstract CharSequence getOrigin(); @@ -87,6 +141,8 @@ public interface Factory { */ public abstract String getW3CTracestate(); + public abstract String getW3CTracestate(SamplingState samplingState); + /** * Stores the original W3C * tracestate header value. @@ -116,6 +172,9 @@ public void updateW3CTracestateFrom(PropagationTags source) { */ public abstract String headerValue(HeaderType headerType, CharSequence lastParentIdOverride); + public abstract String headerValue( + HeaderType headerType, CharSequence lastParentIdOverride, SamplingState samplingState); + /** * Fills a provided tagMap with valid propagated _dd.p.* tags and possibly a new sampling decision * tags _dd.p.dm (root span only) based on the current state, or sets only an error tag if the diff --git a/dd-trace-core/src/main/java/datadog/trace/core/propagation/W3CHttpCodec.java b/dd-trace-core/src/main/java/datadog/trace/core/propagation/W3CHttpCodec.java index 86ea382b482..a6dd7e7b65f 100644 --- a/dd-trace-core/src/main/java/datadog/trace/core/propagation/W3CHttpCodec.java +++ b/dd-trace-core/src/main/java/datadog/trace/core/propagation/W3CHttpCodec.java @@ -23,6 +23,7 @@ import datadog.trace.api.sampling.SamplingMechanism; import datadog.trace.bootstrap.instrumentation.api.TagContext; import datadog.trace.core.DDSpanContext; +import datadog.trace.core.propagation.PropagationTags.SamplingState; import java.util.Map; import java.util.function.Supplier; import org.slf4j.Logger; @@ -65,25 +66,33 @@ public Injector(Map invertedBaggageMapping) { @Override public void inject( final DDSpanContext context, final C carrier, final CarrierSetter setter) { - injectTraceParent(context, carrier, setter); - injectTraceState(context, carrier, setter); + PropagationTags propagationTags = context.getPropagationTags(); + SamplingState samplingState = propagationTags.samplingState(); + injectTraceParent(context, samplingState, carrier, setter); + injectTraceState(context, propagationTags, samplingState, carrier, setter); injectBaggage(context, carrier, setter); } - private void injectTraceParent(DDSpanContext context, C carrier, CarrierSetter setter) { + private void injectTraceParent( + DDSpanContext context, SamplingState samplingState, C carrier, CarrierSetter setter) { String traceparent = W3CTraceParent.from( - context.getTraceId(), context.getSpanId(), context.getSamplingPriority() > 0); + context.getTraceId(), context.getSpanId(), samplingState.getSamplingPriority() > 0); setter.set(carrier, TRACE_PARENT_KEY, traceparent); } - private void injectTraceState(DDSpanContext context, C carrier, CarrierSetter setter) { - PropagationTags propagationTags = context.getPropagationTags(); + private void injectTraceState( + DDSpanContext context, + PropagationTags propagationTags, + SamplingState samplingState, + C carrier, + CarrierSetter setter) { // Supply the injecting span's id for the W3C `p:` as a parameter rather than mutating it into // the (possibly trace-level, shared) tags — keeps transient per-injection identity out of // shared state, so concurrent sibling injects can't race on it. String tracestate = - propagationTags.headerValue(W3C, DDSpanId.toHexStringPadded(context.getSpanId())); + propagationTags.headerValue( + W3C, DDSpanId.toHexStringPadded(context.getSpanId()), samplingState); if (tracestate != null && !tracestate.isEmpty()) { setter.set(carrier, TRACE_STATE_KEY, tracestate); } diff --git a/dd-trace-core/src/main/java/datadog/trace/core/propagation/ptags/DatadogPTagsCodec.java b/dd-trace-core/src/main/java/datadog/trace/core/propagation/ptags/DatadogPTagsCodec.java index 3ac0c7ad712..f8e745ec10c 100644 --- a/dd-trace-core/src/main/java/datadog/trace/core/propagation/ptags/DatadogPTagsCodec.java +++ b/dd-trace-core/src/main/java/datadog/trace/core/propagation/ptags/DatadogPTagsCodec.java @@ -3,6 +3,7 @@ import datadog.logging.RatelimitedLogger; import datadog.trace.api.ProductTraceSource; import datadog.trace.core.propagation.PropagationTags; +import datadog.trace.core.propagation.PropagationTags.SamplingState; import datadog.trace.core.propagation.ptags.PTagsFactory.PTags; import datadog.trace.core.propagation.ptags.TagElement.Encoding; import java.util.ArrayList; @@ -123,14 +124,18 @@ PropagationTags fromHeaderValue(PTagsFactory tagsFactory, String value) { } @Override - protected int estimateHeaderSize(PTags pTags) { - return pTags.getXDatadogTagsSize(); + protected int estimateHeaderSize( + PTags pTags, CharSequence lastParentIdOverride, SamplingState samplingState) { + return pTags.getXDatadogTagsSize(samplingState); } @Override - protected int appendPrefix(StringBuilder sb, PTags ptags) { - // Calculate the tag size here and return it. Don't do anything else since there is no prefix. - return ptags.getXDatadogTagsSize(); + protected int appendPrefix( + StringBuilder sb, + PTags ptags, + CharSequence lastParentIdOverride, + SamplingState samplingState) { + return ptags.getXDatadogTagsSize(samplingState); } @Override @@ -147,7 +152,7 @@ protected int appendTag(StringBuilder sb, TagElement key, TagElement value, int } @Override - protected int appendSuffix(StringBuilder sb, PTags ptags, int size) { + protected int appendSuffix(StringBuilder sb, PTags ptags, int size, SamplingState samplingState) { return size; } diff --git a/dd-trace-core/src/main/java/datadog/trace/core/propagation/ptags/OtelTraceState.java b/dd-trace-core/src/main/java/datadog/trace/core/propagation/ptags/OtelTraceState.java index dc6456e88ec..b7fc01fdbb1 100644 --- a/dd-trace-core/src/main/java/datadog/trace/core/propagation/ptags/OtelTraceState.java +++ b/dd-trace-core/src/main/java/datadog/trace/core/propagation/ptags/OtelTraceState.java @@ -1,36 +1,440 @@ package datadog.trace.core.propagation.ptags; -final class OtelTraceState { +final class OtelTraceState implements CharSequence { + private static final String RANDOM_VALUE_KEY = "rv:"; + private static final String THRESHOLD_KEY = "th:"; + private static final long HASH_MULTIPLIER = 1_111_111_111_111_111_111L; + private static final long MAX_56_BIT_VALUE = 0x00ff_ffff_ffff_ffffL; + private static final double TWO_TO_56 = 0x1.0p56; + private static final int MAX_VALUE_LENGTH = 256; + private final CharSequence value; - private final int originalPosition; + private final CharSequence fields; private final int originalSize; + private final long randomValue; + private final long threshold; + private final int randomValueStart; + private final int randomValueEnd; + private final int thresholdStart; + private final int thresholdEnd; + private final boolean includeRandomValue; + private final boolean includeThreshold; + private final boolean inheritedRandomValue; + private volatile String materializedValue; - private OtelTraceState(CharSequence value, int originalPosition, int originalSize) { + private OtelTraceState( + CharSequence value, + CharSequence fields, + int originalSize, + long randomValue, + long threshold, + int randomValueStart, + int randomValueEnd, + int thresholdStart, + int thresholdEnd, + boolean includeRandomValue, + boolean includeThreshold, + boolean inheritedRandomValue) { this.value = value; - this.originalPosition = originalPosition; + this.fields = fields; this.originalSize = originalSize; + this.randomValue = randomValue; + this.threshold = threshold; + this.randomValueStart = randomValueStart; + this.randomValueEnd = randomValueEnd; + this.thresholdStart = thresholdStart; + this.thresholdEnd = thresholdEnd; + this.includeRandomValue = includeRandomValue; + this.includeThreshold = includeThreshold; + this.inheritedRandomValue = inheritedRandomValue; } - static OtelTraceState parse(CharSequence raw, int originalPosition, int originalSize) { - if (raw == null || raw.length() == 0) { + static OtelTraceState parse(CharSequence raw, int originalSize) { + if (raw == null || raw.length() == 0 || raw.length() > MAX_VALUE_LENGTH) { return null; } - return new OtelTraceState(raw, originalPosition, originalSize); + + int randomValueStart = -1; + int randomValueEnd = -1; + int thresholdStart = -1; + int thresholdEnd = -1; + boolean randomValueSeen = false; + boolean invalidRandomValue = false; + boolean hasUnknownField = false; + boolean normalized = false; + int start = 0; + while (start <= raw.length()) { + int end = indexOf(raw, ';', start); + if (end < 0) { + end = raw.length(); + } + if (startsWith(raw, start, end, RANDOM_VALUE_KEY)) { + int candidateStart = start + RANDOM_VALUE_KEY.length(); + if (!randomValueSeen) { + randomValueSeen = true; + if (isLowerHex(raw, candidateStart, end, 14, 14)) { + randomValueStart = candidateStart; + randomValueEnd = end; + } else { + invalidRandomValue = true; + normalized = true; + } + } else { + normalized = true; + } + } else if (startsWith(raw, start, end, THRESHOLD_KEY)) { + int candidateStart = start + THRESHOLD_KEY.length(); + if (thresholdStart < 0 && isLowerHex(raw, candidateStart, end, 1, 14)) { + thresholdStart = candidateStart; + thresholdEnd = end; + } else { + normalized = true; + } + } else if (isUnknownField(raw, start, end)) { + hasUnknownField = true; + } else { + normalized = true; + } + if (end == raw.length()) { + break; + } + start = end + 1; + } + + if (invalidRandomValue) { + thresholdStart = -1; + thresholdEnd = -1; + } + + if (randomValueStart < 0 && thresholdStart < 0 && !hasUnknownField) { + return null; + } + if (normalized) { + CharSequence normalizedValue = + normalize(raw, randomValueStart, randomValueEnd, thresholdStart, thresholdEnd); + return parseCanonical(normalizedValue, originalSize); + } + return new OtelTraceState( + raw, + raw, + originalSize, + parseHex(raw, randomValueStart, randomValueEnd), + parseThreshold(raw, thresholdStart, thresholdEnd), + randomValueStart, + randomValueEnd, + thresholdStart, + thresholdEnd, + randomValueStart >= 0, + thresholdStart >= 0, + true); } - CharSequence getValue() { - return value; + private static OtelTraceState parseCanonical(CharSequence value, int originalSize) { + int randomValueStart = -1; + int randomValueEnd = -1; + int thresholdStart = -1; + int thresholdEnd = -1; + int start = 0; + while (start < value.length()) { + int end = indexOf(value, ';', start); + if (end < 0) { + end = value.length(); + } + if (startsWith(value, start, end, RANDOM_VALUE_KEY)) { + randomValueStart = start + RANDOM_VALUE_KEY.length(); + randomValueEnd = end; + } else if (startsWith(value, start, end, THRESHOLD_KEY)) { + thresholdStart = start + THRESHOLD_KEY.length(); + thresholdEnd = end; + } + start = end + 1; + } + return new OtelTraceState( + value, + value, + originalSize, + parseHex(value, randomValueStart, randomValueEnd), + parseThreshold(value, thresholdStart, thresholdEnd), + randomValueStart, + randomValueEnd, + thresholdStart, + thresholdEnd, + randomValueStart >= 0, + thresholdStart >= 0, + true); } - int length() { - return value.length(); + static OtelTraceState fromProbabilityDecision( + long traceIdLowOrderBits, double rate, boolean sampled) { + long hash = traceIdLowOrderBits * HASH_MULTIPLIER; + long randomValue = (~hash) >>> 8; + long threshold = Math.round((1.0 - rate) * TWO_TO_56); + if (threshold > MAX_56_BIT_VALUE) { + threshold = MAX_56_BIT_VALUE; + } + if (sampled && randomValue < threshold) { + randomValue = threshold; + } else if (!sampled && randomValue >= threshold) { + randomValue = threshold == 0 ? 0 : threshold - 1; + } + + return new OtelTraceState( + null, null, 0, randomValue, threshold, -1, -1, -1, -1, true, true, false); + } + + OtelTraceState withoutThreshold() { + if (!includeThreshold) { + return this; + } + return withFields(includeRandomValue, false, inheritedRandomValue); + } + + OtelTraceState forNonProbabilityDecision() { + return withFields(inheritedRandomValue && includeRandomValue, false, inheritedRandomValue); + } + + boolean isConsistentWith(boolean sampled) { + if (!includeRandomValue || !includeThreshold) { + return true; + } + return (randomValue >= threshold) == sampled; } - int getOriginalPosition() { - return originalPosition; + private OtelTraceState withFields( + boolean retainRandomValue, boolean retainThreshold, boolean randomValueIsInherited) { + if (!retainRandomValue && !retainThreshold && !hasUnknownFields()) { + return null; + } + return new OtelTraceState( + null, + fields, + 0, + randomValue, + threshold, + randomValueStart, + randomValueEnd, + thresholdStart, + thresholdEnd, + retainRandomValue, + retainThreshold, + randomValueIsInherited); } int getOriginalSize() { return originalSize; } + + boolean isMaterialized() { + return materializedValue != null; + } + + @Override + public int length() { + return value == null ? materialize().length() : value.length(); + } + + @Override + public char charAt(int index) { + return value == null ? materialize().charAt(index) : value.charAt(index); + } + + @Override + public CharSequence subSequence(int start, int end) { + return value == null ? materialize().subSequence(start, end) : value.subSequence(start, end); + } + + @Override + public String toString() { + return value == null ? materialize() : value.toString(); + } + + private String materialize() { + String current = materializedValue; + if (current != null) { + return current; + } + StringBuilder result = new StringBuilder(); + if (includeRandomValue) { + appendManagedField( + result, RANDOM_VALUE_KEY, randomValue, randomValueStart, randomValueEnd, false); + } + if (includeThreshold) { + appendManagedField(result, THRESHOLD_KEY, threshold, thresholdStart, thresholdEnd, true); + } + appendUnknownFields(result); + current = result.toString(); + materializedValue = current; + return current; + } + + private void appendManagedField( + StringBuilder result, + String key, + long numericValue, + int sourceStart, + int sourceEnd, + boolean trimTrailingZeros) { + appendSeparator(result); + result.append(key); + if (fields != null && sourceStart >= 0) { + result.append(fields, sourceStart, sourceEnd); + } else { + appendHex(result, numericValue, trimTrailingZeros); + } + } + + private void appendUnknownFields(StringBuilder result) { + if (fields == null) { + return; + } + int start = 0; + while (start < fields.length()) { + int end = indexOf(fields, ';', start); + if (end < 0) { + end = fields.length(); + } + if (!startsWith(fields, start, end, RANDOM_VALUE_KEY) + && !startsWith(fields, start, end, THRESHOLD_KEY)) { + appendSeparator(result); + result.append(fields, start, end); + } + start = end + 1; + } + } + + private boolean hasUnknownFields() { + if (fields == null) { + return false; + } + int start = 0; + while (start < fields.length()) { + int end = indexOf(fields, ';', start); + if (end < 0) { + end = fields.length(); + } + if (!startsWith(fields, start, end, RANDOM_VALUE_KEY) + && !startsWith(fields, start, end, THRESHOLD_KEY)) { + return true; + } + start = end + 1; + } + return false; + } + + private static CharSequence normalize( + CharSequence raw, + int randomValueStart, + int randomValueEnd, + int thresholdStart, + int thresholdEnd) { + StringBuilder result = new StringBuilder(raw.length()); + if (randomValueStart >= 0) { + appendRange(result, RANDOM_VALUE_KEY, raw, randomValueStart, randomValueEnd); + } + if (thresholdStart >= 0) { + appendRange(result, THRESHOLD_KEY, raw, thresholdStart, thresholdEnd); + } + int start = 0; + while (start < raw.length()) { + int end = indexOf(raw, ';', start); + if (end < 0) { + end = raw.length(); + } + if (!startsWith(raw, start, end, RANDOM_VALUE_KEY) + && !startsWith(raw, start, end, THRESHOLD_KEY) + && isUnknownField(raw, start, end)) { + appendSeparator(result); + result.append(raw, start, end); + } + start = end + 1; + } + return result.toString(); + } + + private static void appendRange( + StringBuilder result, String key, CharSequence source, int start, int end) { + appendSeparator(result); + result.append(key).append(source, start, end); + } + + private static void appendSeparator(StringBuilder result) { + if (result.length() > 0) { + result.append(';'); + } + } + + private static void appendHex(StringBuilder result, long value, boolean trimTrailingZeros) { + int digits = 14; + if (trimTrailingZeros) { + long remaining = value; + while (digits > 1 && (remaining & 0xf) == 0) { + remaining >>>= 4; + digits--; + } + } + for (int i = 13; i >= 14 - digits; i--) { + int digit = (int) ((value >>> (i * 4)) & 0xf); + result.append((char) (digit < 10 ? '0' + digit : 'a' + digit - 10)); + } + } + + private static int indexOf(CharSequence value, char target, int start) { + for (int i = start; i < value.length(); i++) { + if (value.charAt(i) == target) { + return i; + } + } + return -1; + } + + private static boolean startsWith(CharSequence value, int start, int end, CharSequence prefix) { + if (end - start < prefix.length()) { + return false; + } + for (int i = 0; i < prefix.length(); i++) { + if (value.charAt(start + i) != prefix.charAt(i)) { + return false; + } + } + return true; + } + + private static boolean isUnknownField(CharSequence value, int start, int end) { + int separator = indexOf(value, ':', start); + return separator > start && separator < end - 1; + } + + private static boolean isLowerHex( + CharSequence value, int start, int end, int minimumLength, int maximumLength) { + int length = end - start; + if (length < minimumLength || length > maximumLength) { + return false; + } + for (int i = start; i < end; i++) { + char c = value.charAt(i); + if (!((c >= '0' && c <= '9') || (c >= 'a' && c <= 'f'))) { + return false; + } + } + return true; + } + + private static long parseHex(CharSequence value, int start, int end) { + if (start < 0) { + return -1; + } + long parsed = 0; + for (int i = start; i < end; i++) { + char c = value.charAt(i); + parsed = (parsed << 4) | (c <= '9' ? c - '0' : c - 'a' + 10); + } + return parsed; + } + + private static long parseThreshold(CharSequence value, int start, int end) { + if (start < 0) { + return -1; + } + return parseHex(value, start, end) << (4 * (14 - (end - start))); + } } diff --git a/dd-trace-core/src/main/java/datadog/trace/core/propagation/ptags/PTagsCodec.java b/dd-trace-core/src/main/java/datadog/trace/core/propagation/ptags/PTagsCodec.java index e2c0658a1d2..f09e0f33171 100644 --- a/dd-trace-core/src/main/java/datadog/trace/core/propagation/ptags/PTagsCodec.java +++ b/dd-trace-core/src/main/java/datadog/trace/core/propagation/ptags/PTagsCodec.java @@ -4,6 +4,7 @@ import datadog.trace.api.ProductTraceSource; import datadog.trace.core.propagation.PropagationTags; +import datadog.trace.core.propagation.PropagationTags.SamplingState; import datadog.trace.core.propagation.ptags.PTagsFactory.PTags; import datadog.trace.core.propagation.ptags.TagElement.Encoding; import java.util.Iterator; @@ -24,22 +25,23 @@ abstract class PTagsCodec { protected static final String PROPAGATION_ERROR_INCONSISTENT_TID = "inconsistent_tid "; protected static final TagKey UPSTREAM_SERVICES_DEPRECATED_TAG = TagKey.from("upstream_services"); - static String headerValue(PTagsCodec codec, PTags ptags) { - return headerValue(codec, ptags, null); - } - - static String headerValue(PTagsCodec codec, PTags ptags, CharSequence lastParentIdOverride) { - int estimate = codec.estimateHeaderSize(ptags); + static String headerValue( + PTagsCodec codec, + PTags ptags, + CharSequence lastParentIdOverride, + SamplingState samplingState) { + int estimate = codec.estimateHeaderSize(ptags, lastParentIdOverride, samplingState); if (estimate == 0) { return ""; } // No encoding validation here because we don't allow arbitrary tag change StringBuilder sb = new StringBuilder(estimate); - int size = codec.appendPrefix(sb, ptags, lastParentIdOverride); + int size = codec.appendPrefix(sb, ptags, lastParentIdOverride, samplingState); if (!ptags.isPropagationTagsDisabled()) { - if (ptags.getDecisionMakerTagValue() != null) { - size = codec.appendTag(sb, DECISION_MAKER_TAG, ptags.getDecisionMakerTagValue(), size); + TagValue decisionMakerTagValue = ptags.getDecisionMakerTagValue(samplingState); + if (decisionMakerTagValue != null) { + size = codec.appendTag(sb, DECISION_MAKER_TAG, decisionMakerTagValue, size); } if (ptags.getTraceIdHighOrderBitsHexTagValue() != null) { size = codec.appendTag(sb, TRACE_ID_TAG, ptags.getTraceIdHighOrderBitsHexTagValue(), size); @@ -55,10 +57,9 @@ static String headerValue(PTagsCodec codec, PTags ptags, CharSequence lastParent if (ptags.getDebugPropagation() != null) { size = codec.appendTag(sb, DEBUG_TAG, TagValue.from(ptags.getDebugPropagation()), size); } - if (ptags.getKnuthSamplingRateTagValue() != null) { - size = - codec.appendTag( - sb, KNUTH_SAMPLING_RATE_TAG, ptags.getKnuthSamplingRateTagValue(), size); + TagValue knuthSamplingRateTagValue = ptags.getKnuthSamplingRateTagValue(samplingState); + if (knuthSamplingRateTagValue != null) { + size = codec.appendTag(sb, KNUTH_SAMPLING_RATE_TAG, knuthSamplingRateTagValue, size); } if (ptags.getOrgPropagationMarkerTagValue() != null) { size = @@ -72,7 +73,7 @@ static String headerValue(PTagsCodec codec, PTags ptags, CharSequence lastParent size = codec.appendTag(sb, tagKey, tagValue, size); } } - size = codec.appendSuffix(sb, ptags, size); + size = codec.appendSuffix(sb, ptags, size, samplingState); if (codec.isTooLarge(sb, size)) { return null; } else { @@ -81,7 +82,8 @@ static String headerValue(PTagsCodec codec, PTags ptags, CharSequence lastParent } static void fillTagMap(PTags propagationTags, Map tagMap) { - int newSize = propagationTags.getXDatadogTagsSize(); + SamplingState samplingState = propagationTags.samplingState(); + int newSize = propagationTags.getXDatadogTagsSize(samplingState); if (newSize > propagationTags.getxDatadogTagsLimit()) { // Outgoing x-datadog-tags value length exceeds the configured limit @@ -103,10 +105,11 @@ static void fillTagMap(PTags propagationTags, Map tagMap) { tagKey.forType(Encoding.DATADOG).toString(), tagValue.forType(Encoding.DATADOG).toString()); } - if (propagationTags.getDecisionMakerTagValue() != null) { + TagValue decisionMakerTagValue = propagationTags.getDecisionMakerTagValue(samplingState); + if (decisionMakerTagValue != null) { tagMap.put( DECISION_MAKER_TAG.forType(Encoding.DATADOG).toString(), - propagationTags.getDecisionMakerTagValue().forType(Encoding.DATADOG).toString()); + decisionMakerTagValue.forType(Encoding.DATADOG).toString()); } if (propagationTags.getTraceSource() != ProductTraceSource.UNSET) { tagMap.put( @@ -119,10 +122,12 @@ static void fillTagMap(PTags propagationTags, Map tagMap) { tagMap.put( DEBUG_TAG.forType(Encoding.DATADOG).toString(), propagationTags.getDebugPropagation()); } - if (propagationTags.getKnuthSamplingRateTagValue() != null) { + TagValue knuthSamplingRateTagValue = + propagationTags.getKnuthSamplingRateTagValue(samplingState); + if (knuthSamplingRateTagValue != null) { tagMap.put( KNUTH_SAMPLING_RATE_TAG.forType(Encoding.DATADOG).toString(), - propagationTags.getKnuthSamplingRateTagValue().forType(Encoding.DATADOG).toString()); + knuthSamplingRateTagValue.forType(Encoding.DATADOG).toString()); } if (propagationTags.getOrgPropagationMarkerTagValue() != null) { tagMap.put( @@ -174,21 +179,19 @@ static int calcXDatadogTagsSize(int size, TagKey tagKey, TagValue tagValue) { abstract PropagationTags fromHeaderValue(PTagsFactory tagsFactory, String value); - protected abstract int estimateHeaderSize(PTags pTags); + protected abstract int estimateHeaderSize( + PTags pTags, CharSequence lastParentIdOverride, SamplingState samplingState); - protected abstract int appendPrefix(StringBuilder sb, PTags ptags); - - /** - * Encode the prefix, using {@code lastParentIdOverride} for the W3C {@code p:} when non-null - * (inject-time). Codecs without a last-parent-id (e.g. Datadog) ignore the override. - */ - protected int appendPrefix(StringBuilder sb, PTags ptags, CharSequence lastParentIdOverride) { - return appendPrefix(sb, ptags); - } + protected abstract int appendPrefix( + StringBuilder sb, + PTags ptags, + CharSequence lastParentIdOverride, + SamplingState samplingState); protected abstract int appendTag(StringBuilder sb, TagElement key, TagElement value, int size); - protected abstract int appendSuffix(StringBuilder sb, PTags ptags, int size); + protected abstract int appendSuffix( + StringBuilder sb, PTags ptags, int size, SamplingState samplingState); protected abstract boolean isTooLarge(StringBuilder sb, int size); diff --git a/dd-trace-core/src/main/java/datadog/trace/core/propagation/ptags/PTagsFactory.java b/dd-trace-core/src/main/java/datadog/trace/core/propagation/ptags/PTagsFactory.java index 84f3c700269..c1758c60dce 100644 --- a/dd-trace-core/src/main/java/datadog/trace/core/propagation/ptags/PTagsFactory.java +++ b/dd-trace-core/src/main/java/datadog/trace/core/propagation/ptags/PTagsFactory.java @@ -87,19 +87,17 @@ PropagationTags createInvalid(String error) { static class PTags extends PropagationTags { private static final String EMPTY = ""; + private static final SamplingState EMPTY_SAMPLING_STATE = + new SamplingState(PrioritySampling.UNSET, null, null, null, null); protected final PTagsFactory factory; // tags that don't require any modifications and propagated as-is private final List tagPairs; - @SuppressFBWarnings( - value = "AT_STALE_THREAD_WRITE_OF_PRIMITIVE", - justification = "This field is never accessed concurrently") - private boolean canChangeDecisionMaker; + private final Object samplingStateLock = new Object(); - // extracted decision maker tag for easier updates - private volatile TagValue decisionMakerTagValue; + private boolean canChangeDecisionMaker; private static final AtomicIntegerFieldUpdater TRACE_SOURCE_UPDATER = AtomicIntegerFieldUpdater.newUpdater(PTags.class, "traceSource"); @@ -107,12 +105,9 @@ static class PTags extends PropagationTags { private volatile int traceSource; private volatile String debugPropagation; - private volatile double knuthSamplingRate = Double.NaN; - private volatile TagValue knuthSamplingRateTagValue; - private volatile TagValue orgPropagationMarkerTagValue; - private volatile OtelTraceState otelTraceState; + private volatile SamplingState samplingState; // Static cache for the most-recently-seen rate → TagValue. In steady state a service uses one // rate, so this eliminates the char[] + String allocation on every new PTags instance. @@ -120,12 +115,12 @@ static class PTags extends PropagationTags { private static volatile double cachedKsrRate = Double.NaN; private static volatile TagValue cachedKsrTagValue; - // xDatadogTagsSize of the tagPairs, does not include the decision maker tag - private volatile int xDatadogTagsSize = -1; + private volatile SizeCacheEntry xDatadogTagsSizeCache; - private volatile int samplingPriority; private volatile CharSequence origin; - private volatile String[] headerCache = null; + private volatile HeaderCacheEntry datadogHeaderCache; + private volatile HeaderCacheEntry w3cHeaderCache; + private volatile TracestateCacheEntry tracestateCache; /** The high-order 64 bits of the trace id. */ private volatile long traceIdHighOrderBits; @@ -187,9 +182,8 @@ static class PTags extends PropagationTags { this.factory = factory; this.tagPairs = tagPairs; this.canChangeDecisionMaker = decisionMakerTagValue == null; - this.decisionMakerTagValue = decisionMakerTagValue; this.traceSource = traceSource; - this.samplingPriority = samplingPriority; + this.samplingState = initialSamplingState(samplingPriority, decisionMakerTagValue); this.origin = origin; this.lastParentId = lastParentId; this.orgPropagationMarkerTagValue = orgPropagationMarkerTagValue; @@ -221,24 +215,120 @@ static PTags withError(PTagsFactory factory, String error) { @Override public void updateTraceSamplingPriority(int samplingPriority, int samplingMechanism) { - if (samplingPriority != PrioritySampling.UNSET && canChangeDecisionMaker - || samplingMechanism == SamplingMechanism.EXTERNAL_OVERRIDE) { - doUpdateTraceSamplingPriority(samplingPriority, samplingMechanism); + synchronized (samplingStateLock) { + if (samplingPriority != PrioritySampling.UNSET && canChangeDecisionMaker + || samplingMechanism == SamplingMechanism.EXTERNAL_OVERRIDE) { + OtelTraceState nextOtelTraceState = getOtelTraceState(); + if (nextOtelTraceState != null) { + if (samplingMechanism == SamplingMechanism.EXTERNAL_OVERRIDE + && !nextOtelTraceState.isConsistentWith(samplingPriority > 0)) { + nextOtelTraceState = nextOtelTraceState.withoutThreshold(); + } else if (samplingMechanism != SamplingMechanism.UNKNOWN + && samplingMechanism != SamplingMechanism.EXTERNAL_OVERRIDE) { + nextOtelTraceState = nextOtelTraceState.forNonProbabilityDecision(); + } + } + installSamplingState(samplingPriority, samplingMechanism, nextOtelTraceState); + } } } @Override - public void forceKeep(int samplingMechanism) { - doUpdateTraceSamplingPriority(PrioritySampling.USER_KEEP, samplingMechanism); + public boolean tryUpdateTraceSamplingPriority( + int samplingPriority, int samplingMechanism, boolean allowOverride) { + synchronized (samplingStateLock) { + if (samplingPriority == PrioritySampling.UNSET) { + return false; + } + SamplingState current = samplingState; + if (!allowOverride && current.getSamplingPriority() != PrioritySampling.UNSET) { + return false; + } + OtelTraceState nextOtelTraceState = getOtelTraceState(); + if (nextOtelTraceState != null) { + if ((samplingMechanism == SamplingMechanism.EXTERNAL_OVERRIDE + || samplingMechanism == SamplingMechanism.UNKNOWN) + && !nextOtelTraceState.isConsistentWith(samplingPriority > 0)) { + nextOtelTraceState = nextOtelTraceState.withoutThreshold(); + } else if (samplingMechanism != SamplingMechanism.UNKNOWN) { + nextOtelTraceState = nextOtelTraceState.forNonProbabilityDecision(); + } + } + installSamplingState( + samplingPriority, + samplingMechanism, + nextOtelTraceState, + getKnuthSamplingRateTagValue(), + canChangeDecisionMaker || samplingMechanism == SamplingMechanism.EXTERNAL_OVERRIDE); + return true; + } } - private void doUpdateTraceSamplingPriority(int samplingPriority, int samplingMechanism) { - if (this.samplingPriority != samplingPriority) { - // This should invalidate any cached w3c header - clearCachedHeader(W3C); + @Override + public boolean tryUpdateProbabilitySamplingDecision( + int samplingPriority, + int samplingMechanism, + double sampleRate, + boolean probabilitySamplingResult, + long traceIdLowOrderBits, + boolean allowOverride) { + synchronized (samplingStateLock) { + SamplingState current = samplingState; + if (!allowOverride && current.getSamplingPriority() != PrioritySampling.UNSET) { + return false; + } + OtelTraceState nextOtelTraceState = getOtelTraceState(); + if (nextOtelTraceState == null) { + boolean limiterDemotion = probabilitySamplingResult && samplingPriority <= 0; + if (!limiterDemotion) { + nextOtelTraceState = + OtelTraceState.fromProbabilityDecision( + traceIdLowOrderBits, sampleRate, probabilitySamplingResult); + } + } else if (probabilitySamplingResult && samplingPriority <= 0) { + nextOtelTraceState = nextOtelTraceState.withoutThreshold(); + } + TagValue nextKnuthSamplingRate = knuthSamplingRateTagValue(sampleRate); + installSamplingState( + samplingPriority, + samplingMechanism, + nextOtelTraceState, + nextKnuthSamplingRate, + canChangeDecisionMaker || samplingMechanism == SamplingMechanism.EXTERNAL_OVERRIDE); + return true; + } + } + + @Override + public void forceKeep(int samplingMechanism) { + synchronized (samplingStateLock) { + OtelTraceState nextOtelTraceState = getOtelTraceState(); + if (nextOtelTraceState != null) { + nextOtelTraceState = nextOtelTraceState.forNonProbabilityDecision(); + } + installSamplingState(PrioritySampling.USER_KEEP, samplingMechanism, nextOtelTraceState); } - this.samplingPriority = samplingPriority; - if (samplingPriority > 0) { + } + + private void installSamplingState( + int samplingPriority, int samplingMechanism, OtelTraceState nextOtelTraceState) { + installSamplingState( + samplingPriority, + samplingMechanism, + nextOtelTraceState, + getKnuthSamplingRateTagValue(), + true); + } + + private void installSamplingState( + int samplingPriority, + int samplingMechanism, + OtelTraceState nextOtelTraceState, + TagValue nextKnuthSamplingRateTagValue, + boolean updateDecisionMaker) { + clearCachedHeader(W3C); + TagValue nextDecisionMakerTagValue = getDecisionMakerTagValue(); + if (updateDecisionMaker && samplingPriority > 0) { // TODO should try to keep the old sampling mechanism if we override the value? if (samplingMechanism == SamplingMechanism.EXTERNAL_OVERRIDE) { // There is no specific value for the EXTERNAL_OVERRIDE, so say that it's the DEFAULT @@ -248,22 +338,51 @@ private void doUpdateTraceSamplingPriority(int samplingPriority, int samplingMec // format if (samplingMechanism >= 0) { TagValue newDM = TagValue.from("-" + samplingMechanism); - if (!newDM.equals(decisionMakerTagValue)) { + if (!newDM.equals(nextDecisionMakerTagValue)) { // This should invalidate any cached w3c and datadog header clearCachedHeader(DATADOG); clearCachedHeader(W3C); } - decisionMakerTagValue = newDM; + nextDecisionMakerTagValue = newDM; } - } else { + } else if (updateDecisionMaker) { // Drop the decision maker tag - if (decisionMakerTagValue != null) { + if (nextDecisionMakerTagValue != null) { // This should invalidate any cached w3c and datadog header clearCachedHeader(DATADOG); clearCachedHeader(W3C); } - decisionMakerTagValue = null; + nextDecisionMakerTagValue = null; } + samplingState = + newSamplingState( + samplingPriority, + tracestate, + nextOtelTraceState, + nextDecisionMakerTagValue, + nextKnuthSamplingRateTagValue); + } + + private static SamplingState newSamplingState( + int samplingPriority, + String tracestate, + OtelTraceState otelTraceState, + TagValue decisionMakerTagValue, + TagValue knuthSamplingRateTagValue) { + return new SamplingState( + samplingPriority, + tracestate, + otelTraceState, + decisionMakerTagValue, + knuthSamplingRateTagValue); + } + + private static SamplingState initialSamplingState( + int samplingPriority, TagValue decisionMakerTagValue) { + if (samplingPriority == PrioritySampling.UNSET && decisionMakerTagValue == null) { + return EMPTY_SAMPLING_STATE; + } + return newSamplingState(samplingPriority, null, null, decisionMakerTagValue, null); } @Override @@ -302,26 +421,37 @@ public String getDebugPropagation() { @Override public void updateKnuthSamplingRate(double rate) { - if (Double.compare(knuthSamplingRate, rate) != 0) { - clearCachedHeader(DATADOG); - clearCachedHeader(W3C); - knuthSamplingRate = rate; - if (Double.isNaN(rate)) { - knuthSamplingRateTagValue = null; - } else { - TagValue tv; - if (Double.compare(cachedKsrRate, rate) == 0) { - tv = cachedKsrTagValue; - } else { - tv = TagValue.from(formatKnuthSamplingRate(rate)); - cachedKsrTagValue = tv; - cachedKsrRate = rate; - } - knuthSamplingRateTagValue = tv; + synchronized (samplingStateLock) { + TagValue current = getKnuthSamplingRateTagValue(); + TagValue next = knuthSamplingRateTagValue(rate); + if (!Objects.equals(current, next)) { + clearCachedHeader(DATADOG); + clearCachedHeader(W3C); + SamplingState currentState = samplingState; + samplingState = + newSamplingState( + currentState.getSamplingPriority(), + tracestate, + getOtelTraceState(), + getDecisionMakerTagValue(currentState), + next); } } } + private static TagValue knuthSamplingRateTagValue(double rate) { + if (Double.isNaN(rate)) { + return null; + } + if (Double.compare(cachedKsrRate, rate) == 0) { + return cachedKsrTagValue; + } + TagValue value = TagValue.from(formatKnuthSamplingRate(rate)); + cachedKsrTagValue = value; + cachedKsrRate = rate; + return value; + } + /** * Formats a sampling rate with up to 6 decimal digits of precision and no trailing zeros. * @@ -357,7 +487,11 @@ static String formatKnuthSamplingRate(double rate) { } TagValue getKnuthSamplingRateTagValue() { - return knuthSamplingRateTagValue; + return getKnuthSamplingRateTagValue(samplingState); + } + + TagValue getKnuthSamplingRateTagValue(SamplingState samplingState) { + return asTagValue(samplingState.getKnuthSamplingRate()); } @Override @@ -381,7 +515,12 @@ TagValue getOrgPropagationMarkerTagValue() { @Override public int getSamplingPriority() { - return samplingPriority; + return samplingState.getSamplingPriority(); + } + + @Override + public SamplingState samplingState() { + return samplingState; } @Override @@ -426,14 +565,17 @@ public CharSequence getLastParentId() { @SuppressWarnings("StringEquality") @SuppressFBWarnings("ES_COMPARING_STRINGS_WITH_EQ") public String headerValue(HeaderType headerType) { - String header = getCachedHeader(headerType); + SamplingState currentSamplingState = samplingState; + String header = getCachedHeader(headerType, currentSamplingState); if (header == null) { - header = PTagsCodec.headerValue(factory.getDecoderEncoder(headerType), this); + header = + PTagsCodec.headerValue( + factory.getDecoderEncoder(headerType), this, null, currentSamplingState); if (header != null) { - setCachedHeader(headerType, header); + setCachedHeader(headerType, currentSamplingState, header); } else { // We can still cache the fact that we got back null - setCachedHeader(headerType, EMPTY); + setCachedHeader(headerType, currentSamplingState, EMPTY); } } if (header == EMPTY) { @@ -447,10 +589,22 @@ public String headerValue(HeaderType headerType, CharSequence lastParentIdOverri if (lastParentIdOverride == null) { return headerValue(headerType); } - // Inject-time path: encode fresh with the override; do NOT cache — the W3C `p:` is - // per-injecting-span and these tags may be shared across sibling spans. + SamplingState currentSamplingState = samplingState; String header = - PTagsCodec.headerValue(factory.getDecoderEncoder(headerType), this, lastParentIdOverride); + PTagsCodec.headerValue( + factory.getDecoderEncoder(headerType), + this, + lastParentIdOverride, + currentSamplingState); + return (header == null || header.isEmpty()) ? null : header; + } + + @Override + public String headerValue( + HeaderType headerType, CharSequence lastParentIdOverride, SamplingState samplingState) { + String header = + PTagsCodec.headerValue( + factory.getDecoderEncoder(headerType), this, lastParentIdOverride, samplingState); return (header == null || header.isEmpty()) ? null : header; } @@ -459,31 +613,40 @@ public void fillTagMap(Map tagMap) { PTagsCodec.fillTagMap(this, tagMap); } - private String getCachedHeader(HeaderType headerType) { - String[] cache = headerCache; - if (cache == null) { - return null; - } - return cache[headerType.ordinal()]; + private String getCachedHeader(HeaderType headerType, SamplingState samplingState) { + HeaderCacheEntry cache = headerType == DATADOG ? datadogHeaderCache : w3cHeaderCache; + return cache != null && cache.samplingState == samplingState ? cache.header : null; } - private void setCachedHeader(HeaderType headerType, String header) { - String[] cache = headerCache; - if (cache == null) { - cache = headerCache = new String[HeaderType.getNumValues()]; + private void setCachedHeader( + HeaderType headerType, SamplingState samplingState, String header) { + HeaderCacheEntry entry = new HeaderCacheEntry(samplingState, header); + if (headerType == DATADOG) { + datadogHeaderCache = entry; + } else { + w3cHeaderCache = entry; } - cache[headerType.ordinal()] = header; } private void clearCachedHeader(HeaderType headerType) { if (headerType == DATADOG) { invalidateXDatadogTagsSize(); } - String[] cache = headerCache; - if (cache == null) { - return; + if (headerType == DATADOG) { + datadogHeaderCache = null; + } else { + w3cHeaderCache = null; + } + } + + private static final class HeaderCacheEntry { + private final SamplingState samplingState; + private final String header; + + private HeaderCacheEntry(SamplingState samplingState, String header) { + this.samplingState = samplingState; + this.header = header; } - cache[headerType.ordinal()] = null; } int getxDatadogTagsLimit() { @@ -499,18 +662,20 @@ List getTagPairs() { } private void invalidateXDatadogTagsSize() { - this.xDatadogTagsSize = -1; + xDatadogTagsSizeCache = null; } - int getXDatadogTagsSize() { - int size = xDatadogTagsSize; - if (size == -1) { - size = PTagsCodec.calcXDatadogTagsSize(getTagPairs()); - size = PTagsCodec.calcXDatadogTagsSize(size, DECISION_MAKER_TAG, decisionMakerTagValue); + int getXDatadogTagsSize(SamplingState samplingState) { + SizeCacheEntry cache = xDatadogTagsSizeCache; + if (cache == null || cache.samplingState != samplingState) { + int size = PTagsCodec.calcXDatadogTagsSize(getTagPairs()); + size = + PTagsCodec.calcXDatadogTagsSize( + size, DECISION_MAKER_TAG, getDecisionMakerTagValue(samplingState)); size = PTagsCodec.calcXDatadogTagsSize(size, TRACE_ID_TAG, traceIdHighOrderBitsHexTagValue); size = PTagsCodec.calcXDatadogTagsSize( - size, KNUTH_SAMPLING_RATE_TAG, getKnuthSamplingRateTagValue()); + size, KNUTH_SAMPLING_RATE_TAG, getKnuthSamplingRateTagValue(samplingState)); size = PTagsCodec.calcXDatadogTagsSize( size, ORG_PROPAGATION_MARKER_TAG, getOrgPropagationMarkerTagValue()); @@ -522,9 +687,20 @@ int getXDatadogTagsSize() { TRACE_SOURCE_TAG, TagValue.from(ProductTraceSource.getBitfieldHex(currentProductTraceSource))); } - xDatadogTagsSize = size; + cache = new SizeCacheEntry(samplingState, size); + xDatadogTagsSizeCache = cache; + } + return cache.size; + } + + private static final class SizeCacheEntry { + private final SamplingState samplingState; + private final int size; + + private SizeCacheEntry(SamplingState samplingState, int size) { + this.samplingState = samplingState; + this.size = size; } - return size; } TagValue getTraceIdHighOrderBitsHexTagValue() { @@ -532,7 +708,18 @@ TagValue getTraceIdHighOrderBitsHexTagValue() { } TagValue getDecisionMakerTagValue() { - return decisionMakerTagValue; + return getDecisionMakerTagValue(samplingState); + } + + TagValue getDecisionMakerTagValue(SamplingState samplingState) { + return asTagValue(samplingState.getDecisionMaker()); + } + + private static TagValue asTagValue(CharSequence value) { + if (value == null) { + return null; + } + return value instanceof TagValue ? (TagValue) value : TagValue.from(value); } @Override @@ -540,6 +727,27 @@ public String getW3CTracestate() { return this.tracestate; } + @Override + public String getW3CTracestate(SamplingState samplingState) { + TracestateCacheEntry cache = tracestateCache; + if (cache == null || cache.samplingState != samplingState) { + cache = + new TracestateCacheEntry(samplingState, W3CPTagsCodec.rebuildTracestate(samplingState)); + tracestateCache = cache; + } + return cache.tracestate; + } + + private static final class TracestateCacheEntry { + private final SamplingState samplingState; + private final String tracestate; + + private TracestateCacheEntry(SamplingState samplingState, String tracestate) { + this.samplingState = samplingState; + this.tracestate = tracestate; + } + } + @Override public void updateW3CTracestate(String tracestate) { setW3CTracestate(tracestate, W3CPTagsCodec.extractOtelTraceState(tracestate)); @@ -552,24 +760,51 @@ public void updateW3CTracestateFrom(PropagationTags source) { return; } PTags sourcePTags = (PTags) source; - setW3CTracestate(sourcePTags.tracestate, sourcePTags.getOtelTraceState()); + SamplingState sourceState = sourcePTags.samplingState(); + CharSequence sourceOtelTraceState = sourceState.getOtelTraceState(); + setW3CTracestate( + sourceState.getTracestate(), + sourceOtelTraceState instanceof OtelTraceState + ? (OtelTraceState) sourceOtelTraceState + : W3CPTagsCodec.extractOtelTraceState(sourceState.getTracestate())); } private void setW3CTracestate(String tracestate, OtelTraceState otelTraceState) { - clearCachedHeader(W3C); - this.tracestate = tracestate; - this.otelTraceState = otelTraceState; + synchronized (samplingStateLock) { + clearCachedHeader(W3C); + int samplingPriority = samplingState.getSamplingPriority(); + if (otelTraceState != null + && samplingPriority != PrioritySampling.UNSET + && !otelTraceState.isConsistentWith(samplingPriority > 0)) { + otelTraceState = otelTraceState.withoutThreshold(); + } + this.tracestate = tracestate; + this.samplingState = + newSamplingState( + samplingPriority, + tracestate, + otelTraceState, + getDecisionMakerTagValue(), + getKnuthSamplingRateTagValue()); + } } - OtelTraceState getOtelTraceState() { - return otelTraceState; + private OtelTraceState getOtelTraceState() { + return (OtelTraceState) samplingState.getOtelTraceState(); } void setOtelTraceState(OtelTraceState otelTraceState) { - if (this.otelTraceState != otelTraceState) { - this.otelTraceState = otelTraceState; + if (getOtelTraceState() != otelTraceState) { clearCachedHeader(W3C); } + SamplingState currentState = samplingState; + this.samplingState = + newSamplingState( + currentState.getSamplingPriority(), + tracestate, + otelTraceState, + getDecisionMakerTagValue(currentState), + getKnuthSamplingRateTagValue(currentState)); } String getError() { @@ -578,12 +813,22 @@ String getError() { @Override public void updateAndLockDecisionMaker(PropagationTags source) { - if (source instanceof PTags) { - canChangeDecisionMaker = false; - decisionMakerTagValue = ((PTags) source).getDecisionMakerTagValue(); - if (decisionMakerTagValue != null) { - clearCachedHeader(DATADOG); - clearCachedHeader(W3C); + synchronized (samplingStateLock) { + if (source instanceof PTags) { + canChangeDecisionMaker = false; + TagValue decisionMakerTagValue = ((PTags) source).getDecisionMakerTagValue(); + if (decisionMakerTagValue != null) { + clearCachedHeader(DATADOG); + clearCachedHeader(W3C); + } + SamplingState currentState = samplingState; + samplingState = + newSamplingState( + currentState.getSamplingPriority(), + tracestate, + getOtelTraceState(), + decisionMakerTagValue, + getKnuthSamplingRateTagValue(currentState)); } } } diff --git a/dd-trace-core/src/main/java/datadog/trace/core/propagation/ptags/W3CPTagsCodec.java b/dd-trace-core/src/main/java/datadog/trace/core/propagation/ptags/W3CPTagsCodec.java index 9cccf3fd45c..1ad40b88f2b 100644 --- a/dd-trace-core/src/main/java/datadog/trace/core/propagation/ptags/W3CPTagsCodec.java +++ b/dd-trace-core/src/main/java/datadog/trace/core/propagation/ptags/W3CPTagsCodec.java @@ -7,9 +7,11 @@ import datadog.trace.api.internal.VisibleForTesting; import datadog.trace.api.sampling.PrioritySampling; import datadog.trace.core.propagation.PropagationTags; +import datadog.trace.core.propagation.PropagationTags.SamplingState; import datadog.trace.core.propagation.ptags.PTagsFactory.PTags; import datadog.trace.core.propagation.ptags.TagElement.Encoding; import datadog.trace.util.SubSequence; +import edu.umd.cs.findbugs.annotations.SuppressFBWarnings; import java.util.ArrayList; import java.util.List; import java.util.concurrent.TimeUnit; @@ -54,7 +56,6 @@ PropagationTags fromHeaderValue(PTagsFactory tagsFactory, String value) { int otelMemberStart = -1; int otelMemberValueStart = -1; int otelMemberValueEnd = -1; - int otelMemberPosition = -1; while (memberStart < len) { if (memberIndex == MAX_MEMBER_COUNT) { // TODO should we return one with an error? @@ -87,7 +88,6 @@ PropagationTags fromHeaderValue(PTagsFactory tagsFactory, String value) { otelMemberStart = memberStart; otelMemberValueStart = memberValueStart; otelMemberValueEnd = memberValueEnd; - otelMemberPosition = memberIndex; } memberIndex++; @@ -104,7 +104,6 @@ PropagationTags fromHeaderValue(PTagsFactory tagsFactory, String value) { otelTraceState = OtelTraceState.parse( SubSequence.of(value, otelMemberValueStart, valueEnd), - otelMemberPosition, memberContributionSize(value, firstMemberStart, otelMemberStart, otelMemberValueEnd)); } @@ -236,50 +235,63 @@ PropagationTags fromHeaderValue(PTagsFactory tagsFactory, String value) { } @Override - protected int estimateHeaderSize(PTags pTags) { - int size = EMPTY_SIZE + 1; // 'dd=' and delimiter; - // Yes, this is a bit much, but better safe than sorry - size += pTags.getXDatadogTagsSize(); + @SuppressWarnings("StringEquality") + @SuppressFBWarnings( + value = "ES_COMPARING_STRINGS_WITH_EQ", + justification = + "Identity determines whether the sampling state retains this PTags instance's raw " + + "tracestate, allowing its parsed size metadata to be reused.") + protected int estimateHeaderSize( + PTags pTags, CharSequence lastParentIdOverride, SamplingState samplingState) { + int size = EMPTY_SIZE + 1; + size += pTags.getXDatadogTagsSize(samplingState); if (pTags.getOrigin() != null) { - size += pTags.getOrigin().length() + 3; // 'o:' + delimiter + size += pTags.getOrigin().length() + 3; } - if (pTags.getSamplingPriority() != PrioritySampling.UNSET) { - size += 5; // 's:-?[0-9]' + delimiter + if (samplingState.getSamplingPriority() != PrioritySampling.UNSET) { + size += 5; } + CharSequence lastParent = + lastParentIdOverride != null ? lastParentIdOverride : pTags.getLastParentId(); + if (lastParent != null) { + size += lastParent.length() + 3; + } + String originalTracestate = samplingState.getTracestate(); boolean includesOriginalTracestate = false; - if (pTags instanceof W3CPTags) { + if (originalTracestate != null + && pTags instanceof W3CPTags + && originalTracestate == pTags.tracestate) { W3CPTags w3CPTags = (W3CPTags) pTags; size += w3CPTags.maxUnknownSize; if (w3CPTags.ddMemberStart != -1) { - size += - (w3CPTags.tracestate.length() - (w3CPTags.ddMemberValueEnd - w3CPTags.ddMemberStart)); + size += originalTracestate.length() - (w3CPTags.ddMemberValueEnd - w3CPTags.ddMemberStart); includesOriginalTracestate = true; } - } else if (pTags.tracestate != null) { - // We assume there is no Datadog list-member - size += pTags.tracestate.length(); + } else if (originalTracestate != null) { + size += originalTracestate.length(); includesOriginalTracestate = true; } - OtelTraceState otelTraceState = pTags.getOtelTraceState(); + CharSequence otelTraceState = samplingState.getOtelTraceState(); if (otelTraceState != null) { - size -= includesOriginalTracestate ? otelTraceState.getOriginalSize() : 0; + if (includesOriginalTracestate && otelTraceState instanceof OtelTraceState) { + size -= ((OtelTraceState) otelTraceState).getOriginalSize(); + } size += OTEL_MEMBER_KEY.length() + otelTraceState.length() + 1; } - return size; - } - - @Override - protected int appendPrefix(StringBuilder sb, PTags ptags) { - return appendPrefix(sb, ptags, null); + return Math.min(size, MAX_HEADER_SIZE); } @Override - protected int appendPrefix(StringBuilder sb, PTags ptags, CharSequence lastParentIdOverride) { + protected int appendPrefix( + StringBuilder sb, + PTags ptags, + CharSequence lastParentIdOverride, + SamplingState samplingState) { sb.append(DATADOG_MEMBER_KEY); // Append sampling priority (s) - if (ptags.getSamplingPriority() != PrioritySampling.UNSET) { + if (samplingState.getSamplingPriority() != PrioritySampling.UNSET) { sb.append("s:"); - sb.append(ptags.getSamplingPriority()); + sb.append(samplingState.getSamplingPriority()); } // Append origin (o) CharSequence origin = ptags.getOrigin(); @@ -317,7 +329,7 @@ protected int appendTag(StringBuilder sb, TagElement key, TagElement value, int } @Override - protected int appendSuffix(StringBuilder sb, PTags ptags, int size) { + protected int appendSuffix(StringBuilder sb, PTags ptags, int size, SamplingState samplingState) { // If there is room for appending unknown from W3CPTags if (size < MAX_HEADER_SIZE && ptags instanceof W3CPTags) { W3CPTags w3cPTags = (W3CPTags) ptags; @@ -329,8 +341,12 @@ protected int appendSuffix(StringBuilder sb, PTags ptags, int size) { sb.setLength(0); size = 0; } + if (size == 0 && canForwardRawTracestate(samplingState)) { + sb.append(samplingState.getTracestate().trim()); + return EMPTY_SIZE + 1; + } // Append the managed OTel member and all other non-Datadog list-members - if (appendOtelAndVendorMembers(sb, ptags, size != 0)) { + if (appendOtelAndVendorMembers(sb, samplingState, size != 0)) { // We don't care about the total size in bytes here, but only the fact that we added something // that should be returned size = Math.max(size, EMPTY_SIZE + 1); @@ -338,6 +354,53 @@ protected int appendSuffix(StringBuilder sb, PTags ptags, int size) { return size; } + private static boolean canForwardRawTracestate(SamplingState samplingState) { + String original = samplingState.getTracestate(); + if (original == null || original.isEmpty()) { + return false; + } + String trimmed = original.trim(); + if (trimmed.isEmpty() || findNextMember(trimmed, 0) != 0) { + return false; + } + CharSequence otelTraceState = samplingState.getOtelTraceState(); + int otelMemberCount = 0; + int memberStart = 0; + while (memberStart < trimmed.length()) { + int memberEnd = trimmed.indexOf(MEMBER_SEPARATOR, memberStart); + if (memberEnd < 0) { + memberEnd = trimmed.length(); + } + if (trimmed.startsWith(DATADOG_MEMBER_KEY, memberStart)) { + return false; + } + if (trimmed.startsWith(OTEL_MEMBER_KEY, memberStart)) { + if (++otelMemberCount > 1 || otelTraceState == null) { + return false; + } + int valueStart = memberStart + OTEL_MEMBER_KEY.length(); + int valueEnd = stripTrailingOWC(trimmed, valueStart, memberEnd); + if (!contentEquals(trimmed, valueStart, valueEnd, otelTraceState)) { + return false; + } + } + memberStart = findNextMember(trimmed, memberEnd + 1); + } + return (otelTraceState == null) == (otelMemberCount == 0); + } + + private static boolean contentEquals(String value, int start, int end, CharSequence expected) { + if (end - start != expected.length()) { + return false; + } + for (int i = 0; i < expected.length(); i++) { + if (value.charAt(start + i) != expected.charAt(i)) { + return false; + } + } + return true; + } + @Override protected boolean isTooLarge(StringBuilder sb, int size) { return size > MAX_HEADER_SIZE; @@ -737,62 +800,114 @@ private static int cleanUpAndAppendUnknown(StringBuilder sb, W3CPTags w3CPTags, } private static boolean appendOtelAndVendorMembers( - StringBuilder sb, PTags ptags, boolean hasDatadogMember) { - String original = ptags.tracestate; - OtelTraceState otelTraceState = ptags.getOtelTraceState(); + StringBuilder sb, SamplingState samplingState, boolean hasDatadogMember) { + String original = samplingState.getTracestate(); + CharSequence otelTraceState = samplingState.getOtelTraceState(); int remainingMembers = MAX_MEMBER_COUNT - (hasDatadogMember ? 1 : 0); - int otherMemberPosition = 0; - int originalMemberPosition = 0; - int otelMemberPositionOffset = 0; - int otelMemberOriginalPosition = - otelTraceState == null ? -1 : otelTraceState.getOriginalPosition(); - boolean otelTraceStateAppended = false; boolean memberAppended = false; + boolean preserveOtelPosition = isUnchangedInheritedOtelMember(original, otelTraceState); + if (!preserveOtelPosition && otelTraceState != null && remainingMembers > 0) { + appendMember(sb, OTEL_MEMBER_KEY, otelTraceState); + remainingMembers--; + memberAppended = true; + } int len = original == null ? 0 : original.length(); int memberStart = original == null ? 0 : findNextMember(original, 0); while (memberStart < len && remainingMembers > 0) { - // Look for member end position int memberEnd = original.indexOf(MEMBER_SEPARATOR, memberStart); if (memberEnd < 0) { memberEnd = len; } boolean datadogMember = original.startsWith(DATADOG_MEMBER_KEY, memberStart); - boolean managedMember = datadogMember || original.startsWith(OTEL_MEMBER_KEY, memberStart); - // offset to correct for dd members that were dropped/relocated before ot's original position - if (datadogMember && originalMemberPosition < otelMemberOriginalPosition) { - otelMemberPositionOffset++; - } - if (!managedMember) { - if (otelTraceState != null - && !otelTraceStateAppended - && otelMemberOriginalPosition - otelMemberPositionOffset == otherMemberPosition) { - appendMember(sb, OTEL_MEMBER_KEY, otelTraceState.getValue()); - remainingMembers--; - otelTraceStateAppended = true; - memberAppended = true; - if (remainingMembers == 0) { - break; - } - } + boolean otelMember = original.startsWith(OTEL_MEMBER_KEY, memberStart); + if (!datadogMember && (!otelMember || preserveOtelPosition)) { int end = stripTrailingOWC(original, memberStart, memberEnd); appendMember(sb, original, memberStart, end); remainingMembers--; - otherMemberPosition++; memberAppended = true; } - originalMemberPosition++; memberStart = findNextMember(original, memberEnd + 1); } - if (otelTraceState != null - && !otelTraceStateAppended - && remainingMembers > 0 - && otelMemberOriginalPosition - otelMemberPositionOffset == otherMemberPosition) { - appendMember(sb, OTEL_MEMBER_KEY, otelTraceState.getValue()); - memberAppended = true; - } return memberAppended; } + private static boolean isUnchangedInheritedOtelMember( + String original, CharSequence otelTraceState) { + if (original == null || otelTraceState == null) { + return false; + } + int otelMemberCount = 0; + int memberStart = findNextMember(original, 0); + while (memberStart < original.length()) { + int memberEnd = original.indexOf(MEMBER_SEPARATOR, memberStart); + if (memberEnd < 0) { + memberEnd = original.length(); + } + if (original.startsWith(OTEL_MEMBER_KEY, memberStart)) { + if (++otelMemberCount > 1) { + return false; + } + int valueStart = memberStart + OTEL_MEMBER_KEY.length(); + int valueEnd = stripTrailingOWC(original, valueStart, memberEnd); + if (!contentEquals(original, valueStart, valueEnd, otelTraceState)) { + return false; + } + } + memberStart = findNextMember(original, memberEnd + 1); + } + return otelMemberCount == 1; + } + + public static String rebuildTracestate(SamplingState samplingState) { + String original = samplingState.getTracestate(); + CharSequence otelTraceState = samplingState.getOtelTraceState(); + // TODO Consider a raw passthrough for unchanged state after checking dd= is not duplicated. + StringBuilder result = new StringBuilder(MAX_HEADER_SIZE); + int memberCount = 0; + + if (original != null) { + int memberStart = findNextMember(original, 0); + while (memberStart < original.length()) { + int memberEnd = original.indexOf(MEMBER_SEPARATOR, memberStart); + if (memberEnd < 0) { + memberEnd = original.length(); + } + if (original.startsWith(DATADOG_MEMBER_KEY, memberStart)) { + int end = stripTrailingOWC(original, memberStart, memberEnd); + appendMember(result, original, memberStart, end); + memberCount++; + break; + } + memberStart = findNextMember(original, memberEnd + 1); + } + } + + if (otelTraceState != null && memberCount < MAX_MEMBER_COUNT) { + appendMember(result, OTEL_MEMBER_KEY, otelTraceState); + memberCount++; + } + + if (original != null) { + int memberStart = findNextMember(original, 0); + while (memberStart < original.length() && memberCount < MAX_MEMBER_COUNT) { + int memberEnd = original.indexOf(MEMBER_SEPARATOR, memberStart); + if (memberEnd < 0) { + memberEnd = original.length(); + } + boolean managed = + original.startsWith(DATADOG_MEMBER_KEY, memberStart) + || original.startsWith(OTEL_MEMBER_KEY, memberStart); + if (!managed) { + int end = stripTrailingOWC(original, memberStart, memberEnd); + appendMember(result, original, memberStart, end); + memberCount++; + } + memberStart = findNextMember(original, memberEnd + 1); + } + } + return result.length() == 0 ? null : result.toString(); + } + private static void appendMember(StringBuilder sb, String member, int start, int end) { if (sb.length() != 0) { sb.append(MEMBER_SEPARATOR); @@ -811,7 +926,6 @@ static OtelTraceState extractOtelTraceState(String tracestate) { if (tracestate == null || tracestate.isEmpty()) { return null; } - int memberPosition = 0; int firstMemberStart = findNextMember(tracestate, 0); int memberStart = firstMemberStart; int otelMemberStart = -1; @@ -832,7 +946,6 @@ static OtelTraceState extractOtelTraceState(String tracestate) { otelMemberValueEnd = memberValueEnd; break; } - memberPosition++; memberStart = findNextMember(tracestate, memberValueEnd); } if (otelMemberStart == -1) { @@ -841,7 +954,6 @@ static OtelTraceState extractOtelTraceState(String tracestate) { int valueEnd = stripTrailingOWC(tracestate, otelMemberValueStart, otelMemberValueEnd); return OtelTraceState.parse( SubSequence.of(tracestate, otelMemberValueStart, valueEnd), - memberPosition, memberContributionSize(tracestate, firstMemberStart, otelMemberStart, otelMemberValueEnd)); } @@ -852,8 +964,9 @@ private static int memberContributionSize( return isOnlyMember ? memberSize : memberSize + 1; } + /** Creates tags that preserve unmanaged W3C tracestate members without sampling state. */ static W3CPTags empty(PTagsFactory factory, String original) { - return empty(factory, original, extractOtelTraceState(original)); + return empty(factory, original, null); } private static W3CPTags empty( diff --git a/dd-trace-core/src/test/java/datadog/trace/common/metrics/SimpleSpan.java b/dd-trace-core/src/test/java/datadog/trace/common/metrics/SimpleSpan.java index 41a2a5a0d14..d6463298000 100644 --- a/dd-trace-core/src/test/java/datadog/trace/common/metrics/SimpleSpan.java +++ b/dd-trace-core/src/test/java/datadog/trace/common/metrics/SimpleSpan.java @@ -333,7 +333,11 @@ public SimpleSpan setSamplingPriority(int samplingPriority, int samplingMechanis @Override public SimpleSpan setSamplingPriority( - int samplingPriority, CharSequence rate, double sampleRate, int samplingMechanism) { + int samplingPriority, + CharSequence rate, + double sampleRate, + boolean probabilitySamplingResult, + int samplingMechanism) { return this; } diff --git a/dd-trace-core/src/test/java/datadog/trace/common/sampling/RateByServiceTraceSamplerTest.java b/dd-trace-core/src/test/java/datadog/trace/common/sampling/RateByServiceTraceSamplerTest.java index 6918cf9c31f..cfd76f46734 100644 --- a/dd-trace-core/src/test/java/datadog/trace/common/sampling/RateByServiceTraceSamplerTest.java +++ b/dd-trace-core/src/test/java/datadog/trace/common/sampling/RateByServiceTraceSamplerTest.java @@ -2,6 +2,7 @@ import static datadog.trace.api.sampling.PrioritySampling.SAMPLER_DROP; import static datadog.trace.api.sampling.PrioritySampling.SAMPLER_KEEP; +import static datadog.trace.core.propagation.PropagationTags.HeaderType.W3C; import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertFalse; import static org.junit.jupiter.api.Assertions.assertNull; @@ -216,6 +217,64 @@ void samplingPrioritySetOnContext() { } } + @Test + void defaultFallbackEmitsConsistentProbabilityStateBeforeAgentResponse() { + RateByServiceTraceSampler serviceSampler = new RateByServiceTraceSampler(); + CoreTracer tracer = tracerBuilder().writer(new ListWriter()).build(); + try { + DDSpan span = + (DDSpan) + tracer + .buildSpan("datadog", "fallback") + .withServiceName("spock") + .ignoreActiveSpan() + .start(); + + serviceSampler.setSamplingPriority(span); + + String tracestate = span.spanContext().getPropagationTags().headerValue(W3C); + assertEquals(SAMPLER_KEEP, span.getSamplingPriority()); + assertTrue(tracestate.matches(".*ot=rv:[0-9a-f]{14};th:0.*"), tracestate); + } finally { + tracer.close(); + } + } + + @Test + void emptyAndNullOnlyResponsesKeepDefaultProbabilityFallback() { + RateByServiceTraceSampler serviceSampler = new RateByServiceTraceSampler(); + CoreTracer tracer = tracerBuilder().writer(new ListWriter()).build(); + try { + serviceSampler.onResponse("traces", rateResponse(new String[0][0])); + assertEquals(1.0, serviceSampler.fallbackSampleRate()); + assertDefaultFallbackSampling(serviceSampler, tracer); + + serviceSampler.onResponse("traces", rateResponse("service:,env:", null)); + assertEquals(1.0, serviceSampler.fallbackSampleRate()); + assertDefaultFallbackSampling(serviceSampler, tracer); + } finally { + tracer.close(); + } + } + + private static void assertDefaultFallbackSampling( + RateByServiceTraceSampler serviceSampler, CoreTracer tracer) { + DDSpan span = + (DDSpan) + tracer + .buildSpan("datadog", "fallback") + .withServiceName("spock") + .ignoreActiveSpan() + .start(); + + serviceSampler.setSamplingPriority(span); + + assertEquals(SAMPLER_KEEP, span.getSamplingPriority()); + String tracestate = span.spanContext().getPropagationTags().headerValue(W3C); + assertTrue(tracestate.matches(".*ot=rv:[0-9a-f]{14};th:0.*"), tracestate); + span.finish(); + } + @Test void samplingPrioritySetWhenServiceLater() throws Exception { RateByServiceTraceSampler sampler = new RateByServiceTraceSampler(); 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..3a476c48719 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 @@ -409,7 +409,11 @@ public PojoSpan setSamplingPriority(int samplingPriority, int samplingMechanism) @Override public PojoSpan setSamplingPriority( - int samplingPriority, CharSequence rate, double sampleRate, int samplingMechanism) { + int samplingPriority, + CharSequence rate, + double sampleRate, + boolean probabilitySamplingResult, + int samplingMechanism) { return this; } diff --git a/dd-trace-core/src/test/java/datadog/trace/core/CoreSpanBuilderTest.java b/dd-trace-core/src/test/java/datadog/trace/core/CoreSpanBuilderTest.java index 9d9367551a7..547eeff9531 100644 --- a/dd-trace-core/src/test/java/datadog/trace/core/CoreSpanBuilderTest.java +++ b/dd-trace-core/src/test/java/datadog/trace/core/CoreSpanBuilderTest.java @@ -12,10 +12,14 @@ import static datadog.trace.api.DDTags.THREAD_ID; import static datadog.trace.api.DDTags.THREAD_NAME; import static datadog.trace.api.TracePropagationStyle.DATADOG; +import static datadog.trace.api.sampling.PrioritySampling.USER_KEEP; +import static datadog.trace.api.sampling.SamplingMechanism.LOCAL_USER_RULE; import static datadog.trace.bootstrap.instrumentation.api.AgentTracer.noopSpan; +import static datadog.trace.common.sampling.RuleBasedTraceSampler.SAMPLING_RULE_RATE; import static datadog.trace.test.junit.utils.config.WithConfigExtension.injectSysConfig; import static java.util.concurrent.TimeUnit.MILLISECONDS; import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; import static org.junit.jupiter.api.Assertions.assertNotEquals; import static org.junit.jupiter.api.Assertions.assertNull; import static org.junit.jupiter.api.Assertions.assertTrue; @@ -55,6 +59,14 @@ public class CoreSpanBuilderTest extends DDCoreJavaSpecification { + private static final String INHERITED_RANDOM_VALUE = "ef284ace7a91e1"; + private static final String OTEL_TRACE_STATE = + "dd=s:0,ot=rv:" + INHERITED_RANDOM_VALUE + ";th:e6666666666668"; + private static final String OTEL_MEMBER = "ot="; + private static final String THRESHOLD_0_5 = ";th:8"; + private static final double SAMPLE_RATE_0_5 = 0.5; + private static final String DATADOG_TRACE_STATE = "_dd.p.dm=934086a686-4,_dd.p.anytag=value"; + private ListWriter writer; private CoreTracer tracer; @@ -375,9 +387,7 @@ void buildContextFromExtractedContextWithRestartBehavior() { Collections.emptyMap(), Collections.emptyMap(), null, - PropagationTags.factory() - .fromHeaderValue( - PropagationTags.HeaderType.DATADOG, "_dd.p.dm=934086a686-4,_dd.p.anytag=value"), + propagationTagsWithOtelState(), null, DATADOG); DDSpan span = (DDSpan) tracer.buildSpan("test", "op name").asChildOf(extractedContext).start(); @@ -394,6 +404,16 @@ void buildContextFromExtractedContextWithRestartBehavior() { assertEquals( extractedContext.getPropagationTags().headerValue(PropagationTags.HeaderType.W3C), link.traceState()); + String initialTraceState = + span.spanContext().getPropagationTags().headerValue(PropagationTags.HeaderType.W3C); + assertTrue(initialTraceState == null || !initialTraceState.contains(OTEL_MEMBER)); + + span.setSamplingPriority(USER_KEEP, SAMPLING_RULE_RATE, SAMPLE_RATE_0_5, true, LOCAL_USER_RULE); + + String freshTraceState = + span.spanContext().getPropagationTags().headerValue(PropagationTags.HeaderType.W3C); + assertFalse(freshTraceState.contains(INHERITED_RANDOM_VALUE)); + assertTrue(freshTraceState.contains(THRESHOLD_0_5)); } @Test @@ -409,9 +429,7 @@ void buildContextFromExtractedContextWithIgnoreBehavior() { Collections.emptyMap(), Collections.emptyMap(), null, - PropagationTags.factory() - .fromHeaderValue( - PropagationTags.HeaderType.DATADOG, "_dd.p.dm=934086a686-4,_dd.p.anytag=value"), + propagationTagsWithOtelState(), null, DATADOG); DDSpan span = (DDSpan) tracer.buildSpan("test", "op name").asChildOf(extractedContext).start(); @@ -420,6 +438,24 @@ void buildContextFromExtractedContextWithIgnoreBehavior() { assertNotEquals(extractedContext.getSpanId(), span.getParentId()); assertEquals(PrioritySampling.UNSET, span.samplingPriority()); assertTrue(span.getLinks().isEmpty()); + String initialTraceState = + span.spanContext().getPropagationTags().headerValue(PropagationTags.HeaderType.W3C); + assertTrue(initialTraceState == null || !initialTraceState.contains(OTEL_MEMBER)); + + span.setSamplingPriority(USER_KEEP, SAMPLING_RULE_RATE, SAMPLE_RATE_0_5, true, LOCAL_USER_RULE); + + String freshTraceState = + span.spanContext().getPropagationTags().headerValue(PropagationTags.HeaderType.W3C); + assertFalse(freshTraceState.contains(INHERITED_RANDOM_VALUE)); + assertTrue(freshTraceState.contains(THRESHOLD_0_5)); + } + + private static PropagationTags propagationTagsWithOtelState() { + PropagationTags propagationTags = + PropagationTags.factory() + .fromHeaderValue(PropagationTags.HeaderType.DATADOG, DATADOG_TRACE_STATE); + propagationTags.updateW3CTracestate(OTEL_TRACE_STATE); + return propagationTags; } @Test diff --git a/dd-trace-core/src/test/java/datadog/trace/core/OtelSamplingDecisionTest.java b/dd-trace-core/src/test/java/datadog/trace/core/OtelSamplingDecisionTest.java new file mode 100644 index 00000000000..12e489f21a6 --- /dev/null +++ b/dd-trace-core/src/test/java/datadog/trace/core/OtelSamplingDecisionTest.java @@ -0,0 +1,164 @@ +package datadog.trace.core; + +import static datadog.trace.api.config.TracerConfig.TRACE_RATE_LIMIT; +import static datadog.trace.api.config.TracerConfig.TRACE_SAMPLE_RATE; +import static datadog.trace.api.config.TracerConfig.TRACE_SAMPLING_RULES; +import static datadog.trace.api.sampling.PrioritySampling.USER_DROP; +import static datadog.trace.api.sampling.PrioritySampling.USER_KEEP; +import static datadog.trace.api.sampling.SamplingMechanism.LOCAL_USER_RULE; +import static datadog.trace.common.sampling.RuleBasedTraceSampler.SAMPLING_RULE_RATE; +import static datadog.trace.core.propagation.PropagationTags.HeaderType.W3C; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import datadog.trace.common.sampling.PrioritySampler; +import datadog.trace.common.sampling.RateByServiceTraceSampler; +import datadog.trace.common.sampling.Sampler; +import datadog.trace.common.writer.ListWriter; +import java.util.HashMap; +import java.util.Map; +import java.util.Properties; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.ValueSource; + +class OtelSamplingDecisionTest extends DDCoreJavaSpecification { + private static final String AGENT_RATE_ENDPOINT = "traces"; + private static final String OTEL_MEMBER = "ot="; + private static final String OTEL_RANDOM_VALUE_PREFIX = "ot=rv:"; + private static final String HALF_THRESHOLD = ";th:8"; + private static final String MAX_THRESHOLD = ";th:ffffffffffffff"; + private static final double HALF_RATE = 0.5; + private static final String HALF_RATE_RULE = "[{\"sample_rate\": 0.5}]"; + private static final String FULL_RATE_RULE = "[{\"sample_rate\": 1}]"; + + @Test + void initialAgentRateEstablishesDefaultProbabilityState() { + RateByServiceTraceSampler sampler = new RateByServiceTraceSampler(); + withRootSpan( + span -> { + sampler.setSamplingPriority(span); + + String header = w3cHeader(span); + assertTrue(header.contains(OTEL_RANDOM_VALUE_PREFIX)); + assertTrue(header.contains(";th:0")); + }); + } + + @Test + void loadedAgentRateEstablishesProbabilityState() { + RateByServiceTraceSampler sampler = new RateByServiceTraceSampler(); + sampler.onResponse(AGENT_RATE_ENDPOINT, agentRates(HALF_RATE)); + withRootSpan( + span -> { + sampler.setSamplingPriority(span); + + String header = w3cHeader(span); + assertTrue(header.contains(OTEL_RANDOM_VALUE_PREFIX)); + assertTrue(header.contains(HALF_THRESHOLD)); + }); + } + + @Test + void zeroAgentRateUsesDropConsistentMaximumThreshold() { + RateByServiceTraceSampler sampler = new RateByServiceTraceSampler(); + sampler.onResponse(AGENT_RATE_ENDPOINT, agentRates(0)); + withRootSpan( + span -> { + sampler.setSamplingPriority(span); + + String header = w3cHeader(span); + assertTrue(header.contains(OTEL_RANDOM_VALUE_PREFIX)); + assertFalse(header.contains("ot=rv:ffffffffffffff")); + assertTrue(header.contains(MAX_THRESHOLD)); + }); + } + + @ParameterizedTest + @ValueSource(booleans = {true, false}) + void configuredRulesEstablishProbabilityState(boolean traceRule) { + Properties properties = new Properties(); + properties.setProperty( + traceRule ? TRACE_SAMPLING_RULES : TRACE_SAMPLE_RATE, + traceRule ? HALF_RATE_RULE : String.valueOf(HALF_RATE)); + properties.setProperty(TRACE_RATE_LIMIT, "10000000"); + PrioritySampler sampler = (PrioritySampler) Sampler.Builder.forConfig(properties); + withRootSpan( + span -> { + sampler.setSamplingPriority(span); + + String header = w3cHeader(span); + assertTrue(header.contains(OTEL_RANDOM_VALUE_PREFIX)); + assertTrue(header.contains(HALF_THRESHOLD)); + }); + } + + @Test + void limiterRejectionDoesNotFabricateProbabilityState() { + Properties properties = new Properties(); + properties.setProperty(TRACE_SAMPLING_RULES, FULL_RATE_RULE); + properties.setProperty(TRACE_RATE_LIMIT, "1"); + PrioritySampler sampler = (PrioritySampler) Sampler.Builder.forConfig(properties); + CoreTracer tracer = tracerBuilder().writer(new ListWriter()).build(); + try { + DDSpan allowed = newRootSpan(tracer); + DDSpan rejected = newRootSpan(tracer); + + sampler.setSamplingPriority(allowed); + sampler.setSamplingPriority(rejected); + + assertTrue(w3cHeader(allowed).contains(OTEL_RANDOM_VALUE_PREFIX)); + assertEquals(USER_DROP, rejected.samplingPriority()); + assertFalse(w3cHeader(rejected).contains(OTEL_MEMBER)); + } finally { + tracer.close(); + } + } + + @Test + void manualOverrideRemovesLocallyGeneratedProbabilityState() { + withRootSpan( + span -> { + span.setSamplingPriority(USER_KEEP, SAMPLING_RULE_RATE, HALF_RATE, true, LOCAL_USER_RULE); + assertTrue(w3cHeader(span).contains(OTEL_RANDOM_VALUE_PREFIX)); + + span.spanContext().forceKeep(); + + assertFalse(w3cHeader(span).contains(OTEL_MEMBER)); + }); + } + + private void withRootSpan(java.util.function.Consumer test) { + CoreTracer tracer = tracerBuilder().writer(new ListWriter()).build(); + try { + test.accept(newRootSpan(tracer)); + } finally { + tracer.close(); + } + } + + private static DDSpan newRootSpan(CoreTracer tracer) { + return (DDSpan) + tracer + .buildSpan("datadog", "operation") + .withServiceName("service") + .ignoreActiveSpan() + .start(); + } + + private static String w3cHeader(DDSpan span) { + String header = span.spanContext().getPropagationTags().headerValue(W3C); + assertNotNull(header); + return header; + } + + private static Map> agentRates(double rate) { + Map rates = new HashMap<>(); + rates.put("service:,env:", rate); + Map> response = new HashMap<>(); + response.put("rate_by_service", rates); + return response; + } +} diff --git a/dd-trace-core/src/test/java/datadog/trace/core/otlp/trace/OtlpTraceJsonCollectorTest.java b/dd-trace-core/src/test/java/datadog/trace/core/otlp/trace/OtlpTraceJsonCollectorTest.java index 17ec0c0d17f..6879cc976ef 100644 --- a/dd-trace-core/src/test/java/datadog/trace/core/otlp/trace/OtlpTraceJsonCollectorTest.java +++ b/dd-trace-core/src/test/java/datadog/trace/core/otlp/trace/OtlpTraceJsonCollectorTest.java @@ -118,15 +118,16 @@ void nonErrorSpanHasNoStatusObject() throws IOException { } @Test - void spanTraceStateOmittedWhenNotPropagated() throws IOException { + void spanTraceStateIncludesDefaultProbabilityDecision() throws IOException { DDSpan span = startAndFinish("op.notracestate", "GET /no-tracestate", null); OtlpTraceJsonCollector collector = new OtlpTraceJsonCollector(); collector.addTrace(asList((CoreSpan) span)); Map parsedSpan = onlySpan(collector.collectTraces()); - assertFalse( - parsedSpan.containsKey("traceState"), "no W3C tracestate propagated should be omitted"); + assertTrue( + parsedSpan.get("traceState").toString().matches("ot=rv:[0-9a-f]{14};th:0"), + parsedSpan.get("traceState").toString()); } @Test @@ -150,7 +151,9 @@ void spanTraceStateIncludedWhenPropagated() throws IOException { collector.addTrace(asList((CoreSpan) agentSpan)); Map parsedSpan = onlySpan(collector.collectTraces()); - assertEquals("vendor=state", parsedSpan.get("traceState")); + assertTrue( + parsedSpan.get("traceState").toString().matches("ot=rv:[0-9a-f]{14};th:0,vendor=state"), + parsedSpan.get("traceState").toString()); } @Test @@ -181,6 +184,46 @@ void spanFlagsIncludeSampledBitWhenSampled() throws IOException { assertEquals(SAMPLED_TRACE_FLAG, ((Number) parsedSpan.get("flags")).intValue()); } + @Test + void traceStateAndFlagsStayPairedAcrossSamplingDecisions() throws IOException { + Map localFallback = exportSamplingSpan(localProbabilitySpan(1.0, true)); + assertTrue(localFallback.get("traceState").toString().matches("ot=rv:[0-9a-f]{14};th:0")); + assertEquals(SAMPLED_TRACE_FLAG, ((Number) localFallback.get("flags")).intValue()); + + Map inherited = exportSamplingSpan(inheritedSamplingSpan()); + assertEquals("dd=s:1,ot=rv:ef284ace7a91e1;th:8,vendor=state", inherited.get("traceState")); + assertEquals(SAMPLED_TRACE_FLAG, ((Number) inherited.get("flags")).intValue()); + + Map probabilityDrop = exportSamplingSpan(localProbabilitySpan(0.0, false)); + assertTrue( + probabilityDrop + .get("traceState") + .toString() + .matches("ot=rv:[0-9a-f]{14};th:ffffffffffffff")); + assertFalse(probabilityDrop.containsKey("flags")); + + DDSpan limiterDrop = localSamplingSpan(); + limiterDrop + .spanContext() + .getPropagationTags() + .tryUpdateProbabilitySamplingDecision( + PrioritySampling.SAMPLER_DROP, + SamplingMechanism.AGENT_RATE, + 1.0, + true, + limiterDrop.getTraceId().toLong(), + true); + Map limiter = exportSamplingSpan(limiterDrop); + assertNull(limiter.get("traceState")); + assertFalse(limiter.containsKey("flags")); + + DDSpan nonProbabilityKeep = localProbabilitySpan(0.0, false); + nonProbabilityKeep.spanContext().getPropagationTags().forceKeep(SamplingMechanism.MANUAL); + Map nonProbability = exportSamplingSpan(nonProbabilityKeep); + assertNull(nonProbability.get("traceState")); + assertEquals(SAMPLED_TRACE_FLAG, ((Number) nonProbability.get("flags")).intValue()); + } + @Test void multipleSpansInATraceAreAllWritten() throws IOException { AgentSpan parent = TRACER.startSpan("test", "op.parent"); @@ -283,6 +326,54 @@ private static DDSpan startAndFinish(String operationName, String resourceName, return (DDSpan) agentSpan; } + private static DDSpan localSamplingSpan() { + AgentSpan span = TRACER.startSpan("test", "op.sampling"); + span.setResourceName("op.sampling"); + return (DDSpan) span; + } + + private static DDSpan localProbabilitySpan(double rate, boolean sampled) { + DDSpan span = localSamplingSpan(); + span.spanContext() + .getPropagationTags() + .tryUpdateProbabilitySamplingDecision( + sampled ? PrioritySampling.SAMPLER_KEEP : PrioritySampling.SAMPLER_DROP, + SamplingMechanism.AGENT_RATE, + rate, + sampled, + span.getTraceId().toLong(), + true); + return span; + } + + private static DDSpan inheritedSamplingSpan() { + PropagationTags propagationTags = + PropagationTags.factory() + .fromHeaderValue( + PropagationTags.HeaderType.W3C, "dd=s:1,vendor=state,ot=rv:ef284ace7a91e1;th:8"); + ExtractedContext parent = + new ExtractedContext( + DDTraceId.ONE, + 0L, + PrioritySampling.SAMPLER_KEEP, + null, + propagationTags, + TracePropagationStyle.TRACECONTEXT); + AgentSpan span = TRACER.startSpan("test", "op.inherited", parent); + span.setResourceName("op.inherited"); + return (DDSpan) span; + } + + private static Map exportSamplingSpan(DDSpan span) throws IOException { + if (span.getSamplingPriority() <= 0) { + span.setTag(SPAN_SAMPLING_MECHANISM_TAG, SamplingMechanism.SPAN_SAMPLING_RATE); + } + span.finish(); + OtlpTraceJsonCollector collector = new OtlpTraceJsonCollector(); + collector.addTrace(asList((CoreSpan) span)); + return onlySpan(collector.collectTraces()); + } + @SuppressWarnings("unchecked") private static Map onlySpan(OtlpPayload payload) throws IOException { List> spans = allSpans(payload); diff --git a/dd-trace-core/src/test/java/datadog/trace/core/otlp/trace/OtlpTraceProtoTest.java b/dd-trace-core/src/test/java/datadog/trace/core/otlp/trace/OtlpTraceProtoTest.java index aa9d7c7022b..3bb6e7d6eff 100644 --- a/dd-trace-core/src/test/java/datadog/trace/core/otlp/trace/OtlpTraceProtoTest.java +++ b/dd-trace-core/src/test/java/datadog/trace/core/otlp/trace/OtlpTraceProtoTest.java @@ -6,6 +6,7 @@ import static datadog.trace.bootstrap.instrumentation.api.Tags.SPAN_KIND_INTERNAL; import static datadog.trace.bootstrap.instrumentation.api.Tags.SPAN_KIND_PRODUCER; import static datadog.trace.bootstrap.instrumentation.api.Tags.SPAN_KIND_SERVER; +import static datadog.trace.core.DDSpanContext.SPAN_SAMPLING_MECHANISM_TAG; import static datadog.trace.core.otlp.common.OtlpTraceFlags.SAMPLED_TRACE_FLAG; import static java.util.Arrays.asList; import static java.util.Arrays.copyOfRange; @@ -632,6 +633,42 @@ void testCollectMultipleTraces() throws IOException { "payload must contain spans with all three distinct trace IDs"); } + @Test + void traceStateAndFlagsStayPairedAcrossSamplingDecisions() throws IOException { + EncodedSamplingState localFallback = exportSamplingState(localProbabilitySpan(1.0, true)); + assertTrue(localFallback.traceState.matches("ot=rv:[0-9a-f]{14};th:0")); + assertEquals(SAMPLED_TRACE_FLAG, localFallback.flags); + + EncodedSamplingState inherited = exportSamplingState(inheritedSamplingSpan()); + assertEquals("dd=s:1,ot=rv:ef284ace7a91e1;th:8,vendor=state", inherited.traceState); + assertEquals(SAMPLED_TRACE_FLAG, inherited.flags); + + EncodedSamplingState probabilityDrop = exportSamplingState(localProbabilitySpan(0.0, false)); + assertTrue(probabilityDrop.traceState.matches("ot=rv:[0-9a-f]{14};th:ffffffffffffff")); + assertEquals(0, probabilityDrop.flags); + + DDSpan limiterDrop = localSamplingSpan(); + limiterDrop + .spanContext() + .getPropagationTags() + .tryUpdateProbabilitySamplingDecision( + PrioritySampling.SAMPLER_DROP, + SamplingMechanism.AGENT_RATE, + 1.0, + true, + limiterDrop.getTraceId().toLong(), + true); + EncodedSamplingState limiter = exportSamplingState(limiterDrop); + assertNull(limiter.traceState); + assertEquals(0, limiter.flags); + + DDSpan nonProbabilityKeep = localProbabilitySpan(0.0, false); + nonProbabilityKeep.spanContext().getPropagationTags().forceKeep(SamplingMechanism.MANUAL); + EncodedSamplingState nonProbability = exportSamplingState(nonProbabilityKeep); + assertNull(nonProbability.traceState); + assertEquals(SAMPLED_TRACE_FLAG, nonProbability.flags); + } + @Test void poisonedSpanResetsCollectorForNextTrace() { // mid-trace exception (e.g. from a malformed span) must not leave partial state behind @@ -715,6 +752,109 @@ private static List parseSpanNamesFromPayload(OtlpPayload payload) throw return names; } + private static DDSpan localSamplingSpan() { + AgentSpan span = TRACER.startSpan("test", "op.sampling"); + span.setResourceName("op.sampling"); + return (DDSpan) span; + } + + private static DDSpan localProbabilitySpan(double rate, boolean sampled) { + DDSpan span = localSamplingSpan(); + span.spanContext() + .getPropagationTags() + .tryUpdateProbabilitySamplingDecision( + sampled ? PrioritySampling.SAMPLER_KEEP : PrioritySampling.SAMPLER_DROP, + SamplingMechanism.AGENT_RATE, + rate, + sampled, + span.getTraceId().toLong(), + true); + return span; + } + + private static DDSpan inheritedSamplingSpan() { + PropagationTags propagationTags = + PropagationTags.factory() + .fromHeaderValue( + PropagationTags.HeaderType.W3C, "dd=s:1,vendor=state,ot=rv:ef284ace7a91e1;th:8"); + ExtractedContext parent = + new ExtractedContext( + DDTraceId.ONE, + 0L, + PrioritySampling.SAMPLER_KEEP, + null, + propagationTags, + TracePropagationStyle.TRACECONTEXT); + AgentSpan span = TRACER.startSpan("test", "op.inherited", parent); + span.setResourceName("op.inherited"); + return (DDSpan) span; + } + + private static EncodedSamplingState exportSamplingState(DDSpan span) throws IOException { + if (span.getSamplingPriority() <= 0) { + span.setTag(SPAN_SAMPLING_MECHANISM_TAG, SamplingMechanism.SPAN_SAMPLING_RATE); + } + span.finish(); + OtlpTraceProtoCollector collector = new OtlpTraceProtoCollector(); + collector.addTrace(asList((CoreSpan) span)); + return parseOnlySpanSamplingState(collector.collectTraces()); + } + + private static EncodedSamplingState parseOnlySpanSamplingState(OtlpPayload payload) + throws IOException { + CodedInputStream tracesData = CodedInputStream.newInstance(payload.getContent()); + tracesData.readTag(); + CodedInputStream resourceSpans = tracesData.readBytes().newCodedInput(); + CodedInputStream scopeSpans = null; + while (!resourceSpans.isAtEnd()) { + int tag = resourceSpans.readTag(); + if (WireFormat.getTagFieldNumber(tag) == 2) { + scopeSpans = resourceSpans.readBytes().newCodedInput(); + } else { + resourceSpans.skipField(tag); + } + } + assertNotNull(scopeSpans); + + CodedInputStream spanData = null; + while (!scopeSpans.isAtEnd()) { + int tag = scopeSpans.readTag(); + if (WireFormat.getTagFieldNumber(tag) == 2) { + spanData = scopeSpans.readBytes().newCodedInput(); + break; + } + scopeSpans.skipField(tag); + } + assertNotNull(spanData); + + String traceState = null; + int flags = 0; + while (!spanData.isAtEnd()) { + int tag = spanData.readTag(); + switch (WireFormat.getTagFieldNumber(tag)) { + case 3: + traceState = spanData.readString(); + break; + case 16: + flags = spanData.readFixed32(); + break; + default: + spanData.skipField(tag); + } + } + return new EncodedSamplingState(traceState, flags); + } + + private static final class EncodedSamplingState { + private final String traceState; + private final int flags; + + private EncodedSamplingState(String traceState, int flags) { + this.traceState = traceState; + this.flags = flags; + } + } + // ── span construction ───────────────────────────────────────────────────── /** Builds {@link DDSpan} instances from the given specs, collecting them in order. */ diff --git a/dd-trace-core/src/test/java/datadog/trace/core/propagation/OtelTraceStatePropagationTest.java b/dd-trace-core/src/test/java/datadog/trace/core/propagation/OtelTraceStatePropagationTest.java new file mode 100644 index 00000000000..d74ab5f2746 --- /dev/null +++ b/dd-trace-core/src/test/java/datadog/trace/core/propagation/OtelTraceStatePropagationTest.java @@ -0,0 +1,197 @@ +package datadog.trace.core.propagation; + +import static datadog.trace.api.sampling.PrioritySampling.SAMPLER_DROP; +import static datadog.trace.api.sampling.PrioritySampling.SAMPLER_KEEP; +import static datadog.trace.api.sampling.PrioritySampling.USER_KEEP; +import static datadog.trace.api.sampling.SamplingMechanism.AGENT_RATE; +import static datadog.trace.api.sampling.SamplingMechanism.EXTERNAL_OVERRIDE; +import static datadog.trace.api.sampling.SamplingMechanism.MANUAL; +import static datadog.trace.core.propagation.PropagationTags.HeaderType.DATADOG; +import static datadog.trace.core.propagation.PropagationTags.HeaderType.W3C; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertNull; +import static org.junit.jupiter.api.Assertions.assertSame; +import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.junit.jupiter.params.provider.Arguments.arguments; + +import datadog.trace.core.propagation.PropagationTags.SamplingState; +import java.util.stream.Stream; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.Arguments; +import org.junit.jupiter.params.provider.MethodSource; + +class OtelTraceStatePropagationTest { + private static final String RV = "ef284ace7a91e1"; + private static final String TH = "e6666666666668"; + + @Test + void reusesEmptySamplingState() { + PropagationTags.Factory factory = PropagationTags.factory(); + + assertSame(factory.empty().samplingState(), factory.empty().samplingState()); + } + + @ParameterizedTest + @MethodSource("inboundTracestates") + void normalizesAndForwardsOnlyFirstManagedOtelMember(String header, String expected) { + PropagationTags tags = PropagationTags.factory().fromHeaderValue(W3C, header); + + assertEquals(expected, tags.headerValue(W3C)); + } + + static Stream inboundTracestates() { + return Stream.of( + arguments("ot=rv:" + RV + ";th:" + TH, "ot=rv:" + RV + ";th:" + TH), + arguments("ot=rv:" + RV, "ot=rv:" + RV), + arguments("ot=th:" + TH, "ot=th:" + TH), + arguments("ot=future:value", "ot=future:value"), + arguments( + "vendor=state,ot=rv:invalid;th:" + TH + ";future:value", + "ot=future:value,vendor=state"), + arguments("vendor=state,ot=rv:invalid;th:invalid", "vendor=state"), + arguments("ot=rv:EF284ACE7A91E1;th:" + TH, null), + arguments( + "vendor=state,ot=rv:" + RV + ",ot=rv:1234567890abcd,other=state", + "ot=rv:" + RV + ",vendor=state,other=state"), + arguments("dd=s:1,dd=s:0,ot=rv:" + RV, "dd=s:1,ot=rv:" + RV)); + } + + @Test + void rebuildsDuplicateDatadogMembers() { + PropagationTags tags = + PropagationTags.factory().fromHeaderValue(W3C, "dd=s:1,dd=s:0,ot=rv:" + RV); + + assertEquals("dd=s:1,ot=rv:" + RV, tags.getW3CTracestate(tags.samplingState())); + } + + @Test + void publishesProbabilityPriorityAndOtelStateTogether() { + PropagationTags tags = PropagationTags.factory().empty(); + SamplingState before = tags.samplingState(); + + assertTrue( + tags.tryUpdateProbabilitySamplingDecision(SAMPLER_KEEP, AGENT_RATE, 1.0, true, 1L, false)); + SamplingState after = tags.samplingState(); + + assertEquals(SAMPLER_KEEP, after.getSamplingPriority()); + assertEquals("-1", after.getDecisionMaker().toString()); + assertEquals("1", after.getKnuthSamplingRate().toString()); + assertTrue(after.getOtelTraceState().toString().matches("rv:[0-9a-f]{14};th:0")); + assertNull(tags.getW3CTracestate(before)); + String tracestate = tags.getW3CTracestate(after); + assertEquals("ot=" + after.getOtelTraceState(), tracestate); + assertSame(tracestate, tags.getW3CTracestate(after)); + } + + @Test + void rejectedSamplingAttemptCannotReplaceProbabilityState() { + PropagationTags tags = PropagationTags.factory().empty(); + assertTrue( + tags.tryUpdateProbabilitySamplingDecision(SAMPLER_KEEP, AGENT_RATE, 0.5, true, 1L, false)); + SamplingState established = tags.samplingState(); + + assertFalse( + tags.tryUpdateProbabilitySamplingDecision(SAMPLER_DROP, AGENT_RATE, 0.1, false, 2L, false)); + + assertEquals(established, tags.samplingState()); + } + + @Test + void atomicPriorityUpdatePreservesLockedDecisionMaker() { + PropagationTags tags = + PropagationTags.factory().fromHeaderValue(DATADOG, "_dd.p.dm=934086a686-4"); + + assertTrue(tags.tryUpdateTraceSamplingPriority(SAMPLER_KEEP, AGENT_RATE, false)); + + SamplingState state = tags.samplingState(); + assertEquals(SAMPLER_KEEP, state.getSamplingPriority()); + assertEquals("934086a686-4", state.getDecisionMaker().toString()); + } + + @Test + void atomicProbabilityUpdatePreservesLockedDecisionMaker() { + PropagationTags tags = + PropagationTags.factory().fromHeaderValue(DATADOG, "_dd.p.dm=934086a686-4"); + + assertTrue( + tags.tryUpdateProbabilitySamplingDecision(SAMPLER_DROP, AGENT_RATE, 0.5, false, 1L, false)); + + SamplingState state = tags.samplingState(); + assertEquals(SAMPLER_DROP, state.getSamplingPriority()); + assertEquals("934086a686-4", state.getDecisionMaker().toString()); + assertEquals("0.5", state.getKnuthSamplingRate().toString()); + assertTrue(state.getOtelTraceState().toString().matches("rv:[0-9a-f]{14};th:8")); + } + + @Test + void inheritedStateRetainsItsPositionRelativeToVendors() { + PropagationTags tags = + PropagationTags.factory() + .fromHeaderValue(W3C, "vendor=state,ot=rv:ef284ace7a91e1;th:e6666666666668,dd=s:1"); + tags.updateTraceSamplingPriority(SAMPLER_KEEP, EXTERNAL_OVERRIDE); + + assertEquals( + "dd=s:1;t.dm:-0,vendor=state,ot=rv:ef284ace7a91e1;th:e6666666666668", + tags.headerValue(W3C)); + } + + @Test + void preservesFinalUnchangedInheritedOtelMember() { + PropagationTags tags = + PropagationTags.factory().fromHeaderValue(W3C, "dd=s:1,first=value,sec=value,ot=rv:" + RV); + + assertEquals("dd=s:1,first=value,sec=value,ot=rv:" + RV, tags.headerValue(W3C)); + } + + @Test + void compoundConflictRemovesThresholdAndRetainsRandomValue() { + PropagationTags tags = + PropagationTags.factory() + .fromHeaderValue(W3C, "dd=s:1,ot=rv:00000000000001;th:8,vendor=state"); + + tags.updateTraceSamplingPriority(SAMPLER_KEEP, EXTERNAL_OVERRIDE); + + assertEquals("dd=s:1;t.dm:-0,ot=rv:00000000000001,vendor=state", tags.headerValue(W3C)); + } + + @Test + void forceKeepRemovesLocallyGeneratedProbabilityState() { + PropagationTags tags = PropagationTags.factory().empty(); + assertTrue( + tags.tryUpdateProbabilitySamplingDecision(SAMPLER_DROP, AGENT_RATE, 0.0, false, 1L, false)); + + tags.forceKeep(MANUAL); + + assertEquals(USER_KEEP, tags.samplingState().getSamplingPriority()); + assertNull(tags.samplingState().getOtelTraceState()); + } + + @Test + void limiterDemotionDoesNotFabricateState() { + PropagationTags tags = PropagationTags.factory().empty(); + + assertTrue( + tags.tryUpdateProbabilitySamplingDecision(SAMPLER_DROP, AGENT_RATE, 1.0, true, 1L, false)); + + assertNull(tags.samplingState().getOtelTraceState()); + } + + @Test + void generatedManagedMembersDisplaceRightmostVendorAtMemberLimit() { + StringBuilder original = new StringBuilder("v0=state"); + for (int i = 1; i < 31; i++) { + original.append(",v").append(i).append("=state"); + } + PropagationTags tags = PropagationTags.factory().fromHeaderValue(W3C, original.toString()); + + assertTrue( + tags.tryUpdateProbabilitySamplingDecision(SAMPLER_KEEP, AGENT_RATE, 0.5, true, 1L, false)); + + String header = tags.headerValue(W3C); + assertEquals(32, header.split(",").length); + assertTrue(header.startsWith("dd=s:1;t.dm:-1;t.ksr:0.5,ot=rv:")); + assertFalse(header.contains("v30=state")); + } +} diff --git a/dd-trace-core/src/test/java/datadog/trace/core/propagation/W3CHttpInjectorTest.java b/dd-trace-core/src/test/java/datadog/trace/core/propagation/W3CHttpInjectorTest.java index 9903867e1c4..9af53d7c4c3 100644 --- a/dd-trace-core/src/test/java/datadog/trace/core/propagation/W3CHttpInjectorTest.java +++ b/dd-trace-core/src/test/java/datadog/trace/core/propagation/W3CHttpInjectorTest.java @@ -1,14 +1,20 @@ package datadog.trace.core.propagation; +import static datadog.trace.api.sampling.PrioritySampling.SAMPLER_DROP; +import static datadog.trace.api.sampling.PrioritySampling.SAMPLER_KEEP; import static datadog.trace.api.sampling.PrioritySampling.UNSET; import static datadog.trace.api.sampling.PrioritySampling.USER_KEEP; +import static datadog.trace.api.sampling.SamplingMechanism.AGENT_RATE; import static datadog.trace.api.sampling.SamplingMechanism.MANUAL; import static datadog.trace.core.propagation.PropagationTags.HeaderType.DATADOG; +import static datadog.trace.core.propagation.PropagationTags.HeaderType.W3C; import static datadog.trace.core.propagation.W3CHttpCodec.OT_BAGGAGE_PREFIX; import static datadog.trace.core.propagation.W3CHttpCodec.TRACE_PARENT_KEY; import static datadog.trace.core.propagation.W3CHttpCodec.TRACE_STATE_KEY; import static java.util.Collections.singletonMap; import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertTrue; import datadog.trace.api.DDSpanId; import datadog.trace.api.DDTraceId; @@ -19,6 +25,9 @@ import datadog.trace.test.junit.utils.converter.TraceIdConverter; import java.util.HashMap; import java.util.Map; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicReference; import org.junit.jupiter.api.Test; import org.junit.jupiter.params.converter.ConvertWith; import org.tabletest.junit.TableTest; @@ -123,6 +132,62 @@ void injectTheDecisionMakerTag() { assertEquals(expected, carrier); } + @Test + void injectUsesSingleSamplingStateAcrossHeaders() throws InterruptedException { + PropagationTags tags = PropagationTags.factory().fromHeaderValue(W3C, "ot=rv:ef284ace7a91e1"); + assertTrue( + tags.tryUpdateProbabilitySamplingDecision(SAMPLER_KEEP, AGENT_RATE, 1.0, true, 1L, true)); + DDSpanContext context = + mockSpanContext( + DDTraceId.from("1"), DDSpanId.from("2"), SAMPLER_KEEP, null, new HashMap<>(), tags); + Map carrier = new HashMap<>(); + CountDownLatch traceparentWritten = new CountDownLatch(1); + CountDownLatch samplingUpdated = new CountDownLatch(1); + AtomicReference failure = new AtomicReference<>(); + Thread updater = + new Thread( + () -> { + try { + assertTrue(traceparentWritten.await(5, TimeUnit.SECONDS)); + assertTrue( + tags.tryUpdateProbabilitySamplingDecision( + SAMPLER_DROP, AGENT_RATE, 0.0, false, 1L, true)); + } catch (Throwable throwable) { + failure.set(throwable); + } finally { + samplingUpdated.countDown(); + } + }); + updater.start(); + + injector.inject( + context, + carrier, + (headers, key, value) -> { + headers.put(key, value); + if (TRACE_PARENT_KEY.equals(key)) { + traceparentWritten.countDown(); + try { + assertTrue(samplingUpdated.await(5, TimeUnit.SECONDS)); + } catch (InterruptedException interrupted) { + Thread.currentThread().interrupt(); + throw new AssertionError(interrupted); + } + } + }); + updater.join(TimeUnit.SECONDS.toMillis(5)); + + assertFalse(updater.isAlive()); + if (failure.get() != null) { + throw new AssertionError(failure.get()); + } + assertTrue(carrier.get(TRACE_PARENT_KEY).endsWith("-01")); + assertEquals( + "dd=s:1;p:0000000000000002;t.dm:-1;t.ksr:1,ot=rv:ef284ace7a91e1", + carrier.get(TRACE_STATE_KEY)); + assertEquals("dd=s:0;t.ksr:0,ot=rv:ef284ace7a91e1", tags.headerValue(W3C)); + } + @Test void updateLastParentIdOnChildSpan() { Map carrier = new HashMap<>(); diff --git a/dd-trace-core/src/test/java/datadog/trace/core/propagation/W3COtelTraceStateContinuationTest.java b/dd-trace-core/src/test/java/datadog/trace/core/propagation/W3COtelTraceStateContinuationTest.java new file mode 100644 index 00000000000..f444c1ded04 --- /dev/null +++ b/dd-trace-core/src/test/java/datadog/trace/core/propagation/W3COtelTraceStateContinuationTest.java @@ -0,0 +1,137 @@ +package datadog.trace.core.propagation; + +import static datadog.trace.api.ConfigDefaults.DEFAULT_TRACE_X_DATADOG_TAGS_MAX_LENGTH; +import static datadog.trace.api.TracePropagationStyle.DATADOG; +import static datadog.trace.api.TracePropagationStyle.TRACECONTEXT; +import static datadog.trace.bootstrap.instrumentation.api.ContextVisitors.stringValuesMap; +import static datadog.trace.core.propagation.HttpCodecTestHelper.headers; +import static datadog.trace.core.propagation.W3CHttpCodec.TRACE_PARENT_KEY; +import static datadog.trace.core.propagation.W3CHttpCodec.TRACE_STATE_KEY; +import static java.util.Arrays.asList; +import static java.util.Collections.emptyMap; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertInstanceOf; +import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.when; + +import datadog.trace.api.Config; +import datadog.trace.bootstrap.instrumentation.api.AgentSpan; +import datadog.trace.core.CoreTracer; +import datadog.trace.core.DDCoreJavaSpecification; +import datadog.trace.core.DDSpanContext; +import java.util.HashMap; +import java.util.LinkedHashSet; +import java.util.Map; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.CsvSource; + +/** Exercises {@code ot=} through real W3C extraction, continuation, and reinjection. */ +class W3COtelTraceStateContinuationTest extends DDCoreJavaSpecification { + private static final String TRACE_PARENT = + "00-00000000000000000000000000000001-123456789abcdef0-01"; + private static final String RANDOM_VALUE = "ef284ace7a91e1"; + private static final String THRESHOLD = "e6666666666668"; + private static final String DD_MEMBER = "dd=s:2;p:123456789abcdef0"; + + private final HttpCodec.Injector injector = W3CHttpCodec.newInjector(emptyMap()); + + @Test + void roundTripsValidOtelState() { + String inbound = DD_MEMBER + ",ot=rv:" + RANDOM_VALUE + ";th:" + THRESHOLD; + + assertTrue(continueTraceAndReinject(inbound).contains("ot=rv:" + RANDOM_VALUE + ";th:")); + } + + @Test + void malformedRandomValueRemovesManagedPair() { + String inbound = DD_MEMBER + ",ot=rv:zz;th:" + THRESHOLD; + String outbound = continueTraceAndReinject(inbound); + + assertFalse(outbound.contains("ot=")); + assertFalse(outbound.contains("th:" + THRESHOLD)); + } + + @ParameterizedTest + @CsvSource({ + "0, 01, ef284ace7a91e1, 00, false", + "2, 00, 00000000000000, 01, false", + "0, 00, 00000000000000, 00, true", + "2, 01, ef284ace7a91e1, 01, true" + }) + void compoundExtractionKeepsFirstPriorityAndReconcilesOtelState( + int datadogPriority, + String inboundFlags, + String randomValue, + String outboundFlags, + boolean thresholdExpected) { + String traceParent = TRACE_PARENT.substring(0, TRACE_PARENT.length() - 2) + inboundFlags; + String inboundTracestate = + DD_MEMBER + ",ot=rv:" + randomValue + ";th:" + THRESHOLD + ",vendor=state"; + Map inboundHeaders = + headers( + DatadogHttpCodec.TRACE_ID_KEY, + "1", + DatadogHttpCodec.SPAN_ID_KEY, + "2", + DatadogHttpCodec.SAMPLING_PRIORITY_KEY, + String.valueOf(datadogPriority), + TRACE_PARENT_KEY, + traceParent, + TRACE_STATE_KEY, + inboundTracestate); + Config config = mock(Config.class); + when(config.getTracePropagationStylesToExtract()) + .thenReturn(new LinkedHashSet<>(asList(DATADOG, TRACECONTEXT))); + when(config.getxDatadogTagsMaxLength()).thenReturn(DEFAULT_TRACE_X_DATADOG_TAGS_MAX_LENGTH); + + CoreTracer tracer = tracerBuilder().build(); + try { + HttpCodec.Extractor extractor = HttpCodec.createExtractor(config, tracer::captureTraceConfig); + ExtractedContext extracted = + assertInstanceOf( + ExtractedContext.class, extractor.extract(inboundHeaders, stringValuesMap())); + assertEquals(inboundTracestate, extracted.getPropagationTags().getW3CTracestate()); + + AgentSpan span = tracer.buildSpan("test", "continued").asChildOf(extracted).start(); + Map outboundHeaders = new HashMap<>(); + injector.inject((DDSpanContext) span.spanContext(), outboundHeaders, Map::put); + span.finish(); + + assertTrue(outboundHeaders.get(TRACE_PARENT_KEY).endsWith("-" + outboundFlags)); + String outboundTracestate = outboundHeaders.get(TRACE_STATE_KEY); + assertTrue(outboundTracestate.startsWith("dd=s:" + datadogPriority)); + assertTrue(outboundTracestate.contains("ot=rv:" + randomValue)); + assertTrue(outboundTracestate.endsWith("vendor=state")); + assertEquals(thresholdExpected, outboundTracestate.contains("th:" + THRESHOLD)); + } finally { + tracer.close(); + } + } + + private String continueTraceAndReinject(String inboundTracestate) { + Map inboundHeaders = + headers(TRACE_PARENT_KEY, TRACE_PARENT, TRACE_STATE_KEY, inboundTracestate); + + CoreTracer tracer = tracerBuilder().build(); + try { + HttpCodec.Extractor extractor = + W3CHttpCodec.newExtractor(Config.get(), tracer::captureTraceConfig); + ExtractedContext extracted = + assertInstanceOf( + ExtractedContext.class, extractor.extract(inboundHeaders, stringValuesMap())); + + AgentSpan span = tracer.buildSpan("test", "continued").asChildOf(extracted).start(); + assertEquals(extracted.getSamplingPriority(), span.getSamplingPriority()); + + Map outboundHeaders = new HashMap<>(); + injector.inject((DDSpanContext) span.spanContext(), outboundHeaders, Map::put); + span.finish(); + return outboundHeaders.get(TRACE_STATE_KEY); + } finally { + tracer.close(); + } + } +} diff --git a/dd-trace-core/src/test/java/datadog/trace/core/propagation/opg/OrgGuardEnforcerTest.java b/dd-trace-core/src/test/java/datadog/trace/core/propagation/opg/OrgGuardEnforcerTest.java index 65d7ee271d4..4eb636264d1 100644 --- a/dd-trace-core/src/test/java/datadog/trace/core/propagation/opg/OrgGuardEnforcerTest.java +++ b/dd-trace-core/src/test/java/datadog/trace/core/propagation/opg/OrgGuardEnforcerTest.java @@ -3,6 +3,7 @@ import static datadog.trace.api.TracePropagationStyle.DATADOG; import static datadog.trace.api.sampling.PrioritySampling.SAMPLER_KEEP; import static datadog.trace.api.sampling.PrioritySampling.UNSET; +import static datadog.trace.api.sampling.SamplingMechanism.AGENT_RATE; import static datadog.trace.api.sampling.SamplingMechanism.MANUAL; import static datadog.trace.core.propagation.PropagationTags.HeaderType.W3C; import static java.util.Collections.emptySet; @@ -152,6 +153,34 @@ void stripPreservesNonDdVendors() { assertTrue(reEncoded.contains("vendor2=def"), "vendor2 missing: " + reEncoded); } + @Test + @DisplayName("strip replaces inherited OTel sampling state after local resampling") + void stripReplacesOtelSamplingStateAfterLocalResampling() { + OrgGuardEnforcer enforcer = enforcer(false, emptySet(), () -> "L"); + PropagationTags tags = + factory.fromHeaderValue( + W3C, "dd=s:1;o:foo;t.opm:upstream-X,ot=rv:00000000000000;th:1,vendor1=abc"); + ExtractedContext ctx = + new ExtractedContext( + DDTraceId.from(123L), + 456L, + SAMPLER_KEEP, + "origin", + tags, + TracePropagationStyle.TRACECONTEXT); + + ExtractedContext stripped = (ExtractedContext) enforcer.enforce(ctx); + assertTrue( + stripped + .getPropagationTags() + .tryUpdateProbabilitySamplingDecision(SAMPLER_KEEP, AGENT_RATE, 1.0, true, 1L, false)); + + String reEncoded = stripped.getPropagationTags().headerValue(W3C); + assertNotNull(reEncoded); + assertTrue( + reEncoded.matches("dd=s:1;t.dm:-1;t.ksr:1,ot=rv:[0-9a-f]{14};th:0,vendor1=abc"), reEncoded); + } + // ---- helpers ---- private OrgGuardEnforcer enforcer( diff --git a/dd-trace-core/src/test/java/datadog/trace/core/propagation/ptags/OtelTraceStateParsingTest.java b/dd-trace-core/src/test/java/datadog/trace/core/propagation/ptags/OtelTraceStateParsingTest.java index 97a6accae40..83e81477ac5 100644 --- a/dd-trace-core/src/test/java/datadog/trace/core/propagation/ptags/OtelTraceStateParsingTest.java +++ b/dd-trace-core/src/test/java/datadog/trace/core/propagation/ptags/OtelTraceStateParsingTest.java @@ -1,43 +1,74 @@ package datadog.trace.core.propagation.ptags; import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; import static org.junit.jupiter.api.Assertions.assertNotNull; import static org.junit.jupiter.api.Assertions.assertNull; -import static org.junit.jupiter.api.Assertions.assertSame; import datadog.trace.util.SubSequence; import org.junit.jupiter.api.Test; class OtelTraceStateParsingTest { private static final String VALUE = "rv:0123456789abcd"; - private static final int INHERITED_POSITION = 2; private static final int ORIGINAL_MEMBER_CONTRIBUTION_SIZE = 21; @Test void ignoresAbsentValues() { - assertNull(OtelTraceState.parse(null, INHERITED_POSITION, ORIGINAL_MEMBER_CONTRIBUTION_SIZE)); - assertNull(OtelTraceState.parse("", INHERITED_POSITION, ORIGINAL_MEMBER_CONTRIBUTION_SIZE)); + assertNull(OtelTraceState.parse(null, ORIGINAL_MEMBER_CONTRIBUTION_SIZE)); + assertNull(OtelTraceState.parse("", ORIGINAL_MEMBER_CONTRIBUTION_SIZE)); } @Test - void retainsValueAndMemberMetadata() { + void retainsValueAndMemberSize() { SubSequence value = SubSequence.of(VALUE, 0, VALUE.length()); - OtelTraceState state = - OtelTraceState.parse(value, INHERITED_POSITION, ORIGINAL_MEMBER_CONTRIBUTION_SIZE); + OtelTraceState state = OtelTraceState.parse(value, ORIGINAL_MEMBER_CONTRIBUTION_SIZE); assertNotNull(state); - assertSame(value, state.getValue()); + assertFalse(state.isMaterialized()); assertEquals(VALUE.length(), state.length()); - assertEquals(INHERITED_POSITION, state.getOriginalPosition()); assertEquals(ORIGINAL_MEMBER_CONTRIBUTION_SIZE, state.getOriginalSize()); } @Test - void extractsOriginalMemberPosition() { + void retainsValidThresholdWithoutRandomValue() { + OtelTraceState state = OtelTraceState.parse("th:8", 0); + + assertEquals("th:8", state.toString()); + } + + @Test + void removesMalformedThresholdAndRetainsValidRandomValue() { + OtelTraceState state = OtelTraceState.parse("rv:0123456789abcd;th:not-hex;x:value", 0); + + assertEquals("rv:0123456789abcd;x:value", state.toString()); + } + + @Test + void malformedRandomValueRemovesManagedPairButRetainsUnknownFields() { + OtelTraceState state = OtelTraceState.parse("rv:invalid;th:8;x:value", 0); + + assertEquals("x:value", state.toString()); + } + + @Test + void malformedFirstRandomValuePreventsRecoveryFromLaterValues() { OtelTraceState state = - W3CPTagsCodec.extractOtelTraceState("first=value,dd=s:1,dd=s:0,ot=" + VALUE); + OtelTraceState.parse("rv:invalid;rv:0123456789abcd;rv:ffffffffffffff;th:8;th:4", 0); - assertNotNull(state); - assertEquals(3, state.getOriginalPosition()); + assertNull(state); + } + + @Test + void keepsFirstValidManagedFields() { + OtelTraceState state = OtelTraceState.parse("rv:0123456789abcd;rv:ffffffffffffff;th:8;th:4", 0); + + assertEquals("rv:0123456789abcd;th:8", state.toString()); + } + + @Test + void rejectsUppercaseAndOverlongManagedFields() { + OtelTraceState state = OtelTraceState.parse("rv:0123456789ABCd;th:123456789abcdef;x:value", 0); + + assertEquals("x:value", state.toString()); } } diff --git a/dd-trace-core/src/test/java/datadog/trace/core/propagation/ptags/OtelTraceStateTest.java b/dd-trace-core/src/test/java/datadog/trace/core/propagation/ptags/OtelTraceStateTest.java new file mode 100644 index 00000000000..0acd0ff21d3 --- /dev/null +++ b/dd-trace-core/src/test/java/datadog/trace/core/propagation/ptags/OtelTraceStateTest.java @@ -0,0 +1,86 @@ +package datadog.trace.core.propagation.ptags; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertNull; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import org.junit.jupiter.api.Test; + +class OtelTraceStateTest { + private static final long DROP_PRECISION_BOUNDARY_TRACE_ID = 5401449561355763072L; + private static final double DROP_PRECISION_BOUNDARY_RATE = 0.05; + private static final String DROP_PRECISION_BOUNDARY_RANDOM_VALUE = "f333333333332f"; + private static final String DROP_PRECISION_BOUNDARY_THRESHOLD = "f333333333333"; + + @Test + void convertsDatadogProbabilityDecision() { + OtelTraceState state = OtelTraceState.fromProbabilityDecision(0xfff972474538efffL, 0.1, true); + + assertFalse(state.isMaterialized()); + assertEquals("rv:ef284ace7a91e1;th:e6666666666668", state.toString()); + assertTrue(state.isMaterialized()); + assertTrue(state.isConsistentWith(true)); + } + + @Test + void serializesThresholds() { + assertThreshold(0.01, "fd70a3d70a3d7"); + assertThreshold(0.1, "e6666666666668"); + assertThreshold(0.2, "ccccccccccccd"); + assertThreshold(0.5, "8"); + assertThreshold(0.99, "028f5c28f5c29"); + assertThreshold(1.0, "0"); + } + + @Test + void rateZeroUsesLargestWireThresholdAndRemainsDropConsistent() { + OtelTraceState state = OtelTraceState.fromProbabilityDecision(0L, 0.0, false); + + assertEquals("rv:fffffffffffffe;th:ffffffffffffff", state.toString()); + assertTrue(state.isConsistentWith(false)); + } + + @Test + void correctsOnlySerializedRandomValueAtKeepBoundary() { + OtelTraceState state = OtelTraceState.fromProbabilityDecision(0x03a93ee8b1999f00L, 0.1, true); + + assertEquals("rv:e6666666666668;th:e6666666666668", state.toString()); + assertTrue(state.isConsistentWith(true)); + } + + @Test + void correctsOnlySerializedRandomValueAtDropBoundary() { + OtelTraceState state = + OtelTraceState.fromProbabilityDecision( + DROP_PRECISION_BOUNDARY_TRACE_ID, DROP_PRECISION_BOUNDARY_RATE, false); + + assertEquals( + "rv:" + DROP_PRECISION_BOUNDARY_RANDOM_VALUE + ";th:" + DROP_PRECISION_BOUNDARY_THRESHOLD, + state.toString()); + assertTrue(state.isConsistentWith(false)); + } + + @Test + void removesLocalRandomnessForNonProbabilityDecision() { + OtelTraceState state = OtelTraceState.fromProbabilityDecision(1L, 1.0, true); + + assertNull(state.forNonProbabilityDecision()); + } + + @Test + void retainsInheritedRandomnessAndUnknownFieldsWithoutThreshold() { + OtelTraceState state = OtelTraceState.parse("rv:0123456789abcd;th:8;x:value", 0); + + OtelTraceState transformed = state.forNonProbabilityDecision(); + + assertEquals("rv:0123456789abcd;x:value", transformed.toString()); + assertTrue(transformed.isConsistentWith(false)); + } + + private static void assertThreshold(double rate, String expectedThreshold) { + OtelTraceState state = OtelTraceState.fromProbabilityDecision(1L, rate, rate > 0.0); + String value = state.toString(); + assertEquals(expectedThreshold, value.substring(value.indexOf(";th:") + 4)); + } +} diff --git a/dd-trace-core/src/traceAgentTest/java/TraceGenerator.java b/dd-trace-core/src/traceAgentTest/java/TraceGenerator.java index 1e1a9c58fe3..d3fc463cbf3 100644 --- a/dd-trace-core/src/traceAgentTest/java/TraceGenerator.java +++ b/dd-trace-core/src/traceAgentTest/java/TraceGenerator.java @@ -359,7 +359,11 @@ public PojoSpan setSamplingPriority(int samplingPriority, int samplingMechanism) @Override public PojoSpan setSamplingPriority( - int samplingPriority, CharSequence rate, double sampleRate, int samplingMechanism) { + int samplingPriority, + CharSequence rate, + double sampleRate, + boolean probabilitySamplingResult, + int samplingMechanism) { return this; } diff --git a/dd-trace-ot/src/ot31CompatibilityTest/java/datadog/opentracing/OT31ApiTest.java b/dd-trace-ot/src/ot31CompatibilityTest/java/datadog/opentracing/OT31ApiTest.java index 5f1e873fb8c..db6dd5ff569 100644 --- a/dd-trace-ot/src/ot31CompatibilityTest/java/datadog/opentracing/OT31ApiTest.java +++ b/dd-trace-ot/src/ot31CompatibilityTest/java/datadog/opentracing/OT31ApiTest.java @@ -127,7 +127,10 @@ void testInjectExtract( + (propagatedPriority > 0 ? ";t.dm:-" + effectiveSamplingMechanism : "") + ";t.tid:" + traceId.toHexStringPadded(32).substring(0, 16) - + (contextPriority == UNSET ? ";t.ksr:1" : ""); + + (contextPriority == UNSET + ? ";t.ksr:1,ot=" + + ddContext.getPropagationTags().samplingState().getOtelTraceState() + : ""); Map expectedTextMap = new HashMap<>(); OTSpanContext otContext = (OTSpanContext) context; diff --git a/dd-trace-ot/src/ot33CompatibilityTest/java/datadog/opentracing/OT33ApiTest.java b/dd-trace-ot/src/ot33CompatibilityTest/java/datadog/opentracing/OT33ApiTest.java index a2d6270bb3a..2b9cfdca200 100644 --- a/dd-trace-ot/src/ot33CompatibilityTest/java/datadog/opentracing/OT33ApiTest.java +++ b/dd-trace-ot/src/ot33CompatibilityTest/java/datadog/opentracing/OT33ApiTest.java @@ -113,7 +113,10 @@ void testInjectExtract( + (propagatedPriority > 0 ? ";t.dm:-" + effectiveSamplingMechanism : "") + ";t.tid:" + traceId.toHexStringPadded(32).substring(0, 16) - + (contextPriority == UNSET ? ";t.ksr:1" : ""); + + (contextPriority == UNSET + ? ";t.ksr:1,ot=" + + ddContext.getPropagationTags().samplingState().getOtelTraceState() + : ""); Map expectedTextMap = new HashMap<>(); expectedTextMap.put("x-datadog-trace-id", context.toTraceId());