Skip to content

Commit db90b17

Browse files
authored
Merge branch 'master' into block-telemetry-2
2 parents 074a0b6 + 79a1710 commit db90b17

35 files changed

Lines changed: 2648 additions & 348 deletions

File tree

Lines changed: 89 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,89 @@
1+
import akka.actor.ActorRef
2+
import akka.actor.ActorSystem
3+
import akka.actor.DeadLetter
4+
import akka.actor.Props
5+
import akka.actor.UntypedActor
6+
import akka.testkit.TestKit
7+
import akka.testkit.TestProbe
8+
import datadog.trace.agent.test.InstrumentationSpecification
9+
import scala.concurrent.duration.Duration
10+
11+
import java.util.concurrent.CountDownLatch
12+
import java.util.concurrent.atomic.AtomicInteger
13+
14+
import static datadog.trace.agent.test.utils.TraceUtils.basicSpan
15+
import static datadog.trace.agent.test.utils.TraceUtils.runUnderTrace
16+
import static java.util.concurrent.TimeUnit.SECONDS
17+
18+
class AkkaDeadLetterTest extends InstrumentationSpecification {
19+
def "release context for #count queued messages discarded on actor termination (dead letter: #deadLetter)"() {
20+
setup:
21+
def system = ActorSystem.create("dead-letter-test")
22+
def entered = new CountDownLatch(1)
23+
def stop = new CountDownLatch(1)
24+
def processed = new AtomicInteger()
25+
def actor = system.actorOf(Props.create(StoppingActor, entered, stop, processed))
26+
def watcher = new TestProbe(system)
27+
watcher.watch(actor)
28+
actor.tell("stop", ActorRef.noSender())
29+
assert entered.await(10, SECONDS)
30+
31+
when:
32+
runUnderTrace("parent") {
33+
count.times {
34+
def message = deadLetter ? new DeadLetter("discarded", system.deadLetters(), actor) : "discarded"
35+
actor.tell(message, ActorRef.noSender())
36+
}
37+
}
38+
39+
then:
40+
TEST_WRITER.empty
41+
42+
when:
43+
stop.countDown()
44+
watcher.expectTerminated(actor, Duration.create(10, SECONDS))
45+
46+
then:
47+
assertTraces(1) {
48+
trace(1) {
49+
basicSpan(it, "parent")
50+
}
51+
}
52+
processed.get() == 0
53+
54+
cleanup:
55+
stop?.countDown()
56+
TestKit.shutdownActorSystem(system, Duration.create(10, SECONDS), true)
57+
58+
where:
59+
count | deadLetter
60+
1 | false
61+
10 | false
62+
1 | true
63+
}
64+
65+
static class StoppingActor extends UntypedActor {
66+
private final CountDownLatch entered
67+
private final CountDownLatch stop
68+
private final AtomicInteger processed
69+
70+
StoppingActor(CountDownLatch entered, CountDownLatch stop, AtomicInteger processed) {
71+
this.entered = entered
72+
this.stop = stop
73+
this.processed = processed
74+
}
75+
76+
@Override
77+
void onReceive(Object message) {
78+
if (message == "stop") {
79+
entered.countDown()
80+
if (!stop.await(10, SECONDS)) {
81+
throw new IllegalStateException("Timed out waiting to stop the actor")
82+
}
83+
context.stop(self)
84+
} else {
85+
processed.incrementAndGet()
86+
}
87+
}
88+
}
89+
}
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,70 @@
1+
package datadog.trace.instrumentation.akka.concurrent;
2+
3+
import static datadog.trace.agent.tooling.bytebuddy.matcher.HierarchyMatchers.implementsInterface;
4+
import static datadog.trace.agent.tooling.bytebuddy.matcher.NameMatchers.nameStartsWith;
5+
import static datadog.trace.agent.tooling.bytebuddy.matcher.NameMatchers.named;
6+
import static datadog.trace.bootstrap.instrumentation.java.concurrent.AdviceUtils.cancelTask;
7+
import static java.util.Collections.singletonMap;
8+
import static net.bytebuddy.matcher.ElementMatchers.isMethod;
9+
import static net.bytebuddy.matcher.ElementMatchers.isPublic;
10+
import static net.bytebuddy.matcher.ElementMatchers.returns;
11+
import static net.bytebuddy.matcher.ElementMatchers.takesArgument;
12+
import static net.bytebuddy.matcher.ElementMatchers.takesArguments;
13+
14+
import akka.dispatch.Envelope;
15+
import com.google.auto.service.AutoService;
16+
import datadog.trace.agent.tooling.Instrumenter;
17+
import datadog.trace.agent.tooling.InstrumenterModule;
18+
import datadog.trace.bootstrap.InstrumentationContext;
19+
import datadog.trace.bootstrap.instrumentation.java.concurrent.State;
20+
import java.util.Map;
21+
import net.bytebuddy.asm.Advice;
22+
import net.bytebuddy.description.type.TypeDescription;
23+
import net.bytebuddy.matcher.ElementMatcher;
24+
25+
@AutoService(InstrumenterModule.class)
26+
public class AkkaDeadLetterQueueInstrumentation extends InstrumenterModule.ContextTracking
27+
implements Instrumenter.ForTypeHierarchy, Instrumenter.HasMethodAdvice {
28+
29+
public AkkaDeadLetterQueueInstrumentation() {
30+
super("akka_actor_receive", "akka_actor", "akka_concurrent", "java_concurrent");
31+
}
32+
33+
@Override
34+
public String hierarchyMarkerType() {
35+
return "akka.dispatch.MessageQueue";
36+
}
37+
38+
@Override
39+
public ElementMatcher<TypeDescription> hierarchyMatcher() {
40+
// Match the dead-letter queue without depending on Scala's anonymous-class numbering.
41+
return nameStartsWith("akka.dispatch.Mailboxes$")
42+
.and(implementsInterface(named("akka.dispatch.MessageQueue")));
43+
}
44+
45+
@Override
46+
public Map<String, String> contextStore() {
47+
return singletonMap("akka.dispatch.Envelope", State.class.getName());
48+
}
49+
50+
@Override
51+
public void methodAdvice(MethodTransformer transformer) {
52+
transformer.applyAdvice(
53+
isMethod()
54+
.and(isPublic())
55+
.and(named("enqueue"))
56+
.and(takesArguments(2))
57+
.and(takesArgument(0, named("akka.actor.ActorRef")))
58+
.and(takesArgument(1, named("akka.dispatch.Envelope")))
59+
.and(returns(void.class)),
60+
getClass().getName() + "$EnqueueAdvice");
61+
}
62+
63+
public static class EnqueueAdvice {
64+
@Advice.OnMethodExit(onThrowable = Throwable.class, suppress = Throwable.class)
65+
public static void afterEnqueue(@Advice.Argument(1) Envelope envelope) {
66+
// Dead-letter publication discards the original envelope without invoking its actor.
67+
cancelTask(InstrumentationContext.get(Envelope.class, State.class), envelope);
68+
}
69+
}
70+
}

‎dd-java-agent/instrumentation/opentelemetry/opentelemetry-0.3/src/test/groovy/OpenTelemetryTest.groovy‎

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -287,6 +287,7 @@ class OpenTelemetryTest extends InstrumentationSpecification {
287287
}
288288
if (contextPriority == UNSET) {
289289
expectedTracestate += ";t.ksr:1"
290+
expectedTracestate += ",ot=rv:[0-9a-f]{14};th:0"
290291
}
291292
if (traceId.toHighOrderLong() != 0) {
292293
expectedDataTags << "_dd.p.tid=" + traceId.toHexStringPadded(32).substring(0, 16)
@@ -299,12 +300,12 @@ class OpenTelemetryTest extends InstrumentationSpecification {
299300
"x-datadog-parent-id" : "$spanId",
300301
"x-datadog-sampling-priority": propagatedPriority.toString(),
301302
"traceparent" : expectedTraceparent,
302-
"tracestate" : expectedTracestate,
303303
]
304304
if (!expectedDataTags.empty) {
305305
expectedTextMap.put("x-datadog-tags", expectedDataTags.join(','))
306306
}
307-
textMap == expectedTextMap
307+
textMap.tracestate ==~ expectedTracestate
308+
textMap.findAll { key, value -> key != "tracestate" } == expectedTextMap
308309

309310
when:
310311
def extractedContext = httpPropagator.extract(context, textMap, new TextMapGetter())

‎dd-java-agent/instrumentation/opentracing/opentracing-0.31/src/test/groovy/OpenTracing31Test.groovy‎

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -295,19 +295,20 @@ class OpenTracing31Test extends InstrumentationSpecification {
295295
}
296296
if (contextPriority == UNSET) {
297297
expectedTracestate += ";t.ksr:1"
298+
expectedTracestate += ",ot=rv:[0-9a-f]{14};th:0"
298299
datadogTags << "_dd.p.ksr=1"
299300
}
300301
def expectedTextMap = [
301302
"x-datadog-trace-id" : "$context.delegate.traceId",
302303
"x-datadog-parent-id" : "$context.delegate.spanId",
303304
"x-datadog-sampling-priority": propagatedPriority.toString(),
304305
"traceparent" : expectedTraceparent,
305-
"tracestate" : expectedTracestate,
306306
]
307307
if (!datadogTags.empty) {
308308
expectedTextMap.put("x-datadog-tags", datadogTags.join(','))
309309
}
310-
textMap == expectedTextMap
310+
textMap.tracestate ==~ expectedTracestate
311+
textMap.findAll { key, value -> key != "tracestate" } == expectedTextMap
311312

312313
when:
313314
def extract = tracer.extract(Format.Builtin.TEXT_MAP, adapter)

‎dd-java-agent/instrumentation/opentracing/opentracing-0.32/src/test/groovy/OpenTracing32Test.groovy‎

Lines changed: 4 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -311,19 +311,20 @@ class OpenTracing32Test extends InstrumentationSpecification {
311311
}
312312
if (contextPriority == UNSET) {
313313
expectedTracestate+= ";t.ksr:1"
314+
expectedTracestate+= ",ot=rv:[0-9a-f]{14};th:0"
314315
datadogTags << "_dd.p.ksr=1"
315316
}
316317
def expectedTextMap = [
317318
"x-datadog-trace-id" : "$context.delegate.traceId",
318319
"x-datadog-parent-id" : "$context.delegate.spanId",
319320
"x-datadog-sampling-priority": propagatedPriority.toString(),
320-
"traceparent" : expectedTraceparent,
321-
"tracestate" : expectedTracestate
321+
"traceparent" : expectedTraceparent
322322
]
323323
if (!datadogTags.empty) {
324324
expectedTextMap.put("x-datadog-tags", datadogTags.join(','))
325325
}
326-
textMap == expectedTextMap
326+
textMap.tracestate ==~ expectedTracestate
327+
textMap.findAll { key, value -> key != "tracestate" } == expectedTextMap
327328

328329
when:
329330
def extract = tracer.extract(Format.Builtin.TEXT_MAP, adapter)

‎dd-trace-core/src/jmh/java/datadog/trace/core/propagation/ptags/KnuthSamplingRateFormatBenchmark.java‎

Lines changed: 12 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,7 @@
11
package datadog.trace.core.propagation.ptags;
22

3+
import static datadog.trace.api.sampling.PrioritySampling.SAMPLER_DROP;
4+
import static datadog.trace.api.sampling.SamplingMechanism.AGENT_RATE;
35
import static java.util.concurrent.TimeUnit.SECONDS;
46

57
import datadog.trace.core.propagation.PropagationTags;
@@ -53,12 +55,14 @@ public class KnuthSamplingRateFormatBenchmark {
5355
@Param({"0.5", "0.1", "0.01", "0.001", "0.0001", "0.123456789", "0.999999"})
5456
double rate;
5557

58+
PropagationTags.Factory factory;
5659
PTagsFactory.PTags ptags;
5760

5861
@Setup(Level.Trial)
5962
public void setUp() {
60-
ptags = (PTagsFactory.PTags) PropagationTags.factory().empty();
61-
ptags.updateKnuthSamplingRate(rate);
63+
factory = PropagationTags.factory();
64+
ptags = (PTagsFactory.PTags) factory.empty();
65+
ptags.tryUpdateProbabilitySamplingDecision(SAMPLER_DROP, AGENT_RATE, rate, false, 1L, false);
6266
}
6367

6468
/** Baseline: old implementation using String.format + substring trimming. */
@@ -82,16 +86,13 @@ public void cachedTagValue(Blackhole bh) {
8286
bh.consume(ptags.getKnuthSamplingRateTagValue());
8387
}
8488

85-
/**
86-
* Models the per-trace allocation cost: resets the instance cache (simulating a new PTags), then
87-
* calls updateKnuthSamplingRate. This is what every trace root pays. With the static cache
88-
* applied, this should also be near-zero allocation after warmup.
89-
*/
89+
/** Models a fresh trace's probability decision, including its propagation tags allocation. */
9090
@Benchmark
91-
public void updateRateFreshTrace(Blackhole bh) {
92-
ptags.updateKnuthSamplingRate(Double.NaN); // reset instance cache, like a new PTags
93-
ptags.updateKnuthSamplingRate(rate);
94-
bh.consume(ptags.getKnuthSamplingRateTagValue());
91+
public void probabilityDecisionFreshTrace(Blackhole bh) {
92+
PTagsFactory.PTags freshTags = (PTagsFactory.PTags) factory.empty();
93+
freshTags.tryUpdateProbabilitySamplingDecision(
94+
SAMPLER_DROP, AGENT_RATE, rate, false, 1L, false);
95+
bh.consume(freshTags.getKnuthSamplingRateTagValue());
9596
}
9697

9798
// ---- old implementation for comparison (%.6f with trailing zero removal) ----

‎dd-trace-core/src/main/java/datadog/trace/common/sampling/RateByServiceTraceSampler.java‎

Lines changed: 7 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -64,19 +64,13 @@ public <T extends CoreSpan<T>> void setSamplingPriority(final T span) {
6464
final RateSamplersByEnvAndService rates = serviceRates;
6565
RateSampler sampler = rates.getSampler(env, serviceName);
6666

67-
if (sampler.sample(span)) {
68-
span.setSamplingPriority(
69-
PrioritySampling.SAMPLER_KEEP,
70-
SAMPLING_AGENT_RATE,
71-
sampler.getSampleRate(),
72-
SamplingMechanism.AGENT_RATE);
73-
} else {
74-
span.setSamplingPriority(
75-
PrioritySampling.SAMPLER_DROP,
76-
SAMPLING_AGENT_RATE,
77-
sampler.getSampleRate(),
78-
SamplingMechanism.AGENT_RATE);
79-
}
67+
boolean sampled = sampler.sample(span);
68+
int samplingPriority = sampled ? PrioritySampling.SAMPLER_KEEP : PrioritySampling.SAMPLER_DROP;
69+
span.setSamplingPriority(
70+
samplingPriority,
71+
SAMPLING_AGENT_RATE,
72+
sampler.getSampleRate(),
73+
SamplingMechanism.AGENT_RATE);
8074
}
8175

8276
private <T extends CoreSpan<T>> String getSpanEnv(final T span) {

‎dd-trace-core/src/main/java/datadog/trace/common/sampling/RuleBasedTraceSampler.java‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -186,7 +186,8 @@ public <T extends CoreSpan<T>> void setSamplingPriority(final T span) {
186186
PrioritySampling.USER_DROP,
187187
SAMPLING_RULE_RATE,
188188
matchedRule.getSampler().getSampleRate(),
189-
matchedRule.getMechanism());
189+
matchedRule.getMechanism(),
190+
true);
190191
}
191192
span.setMetric(SAMPLING_LIMIT_RATE, rateLimit);
192193
} else {

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

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -124,6 +124,15 @@ default void processTagsAndBaggageWithStructuredLinks(
124124
T setSamplingPriority(
125125
int samplingPriority, CharSequence rate, double sampleRate, int samplingMechanism);
126126

127+
default T setSamplingPriority(
128+
int samplingPriority,
129+
CharSequence rate,
130+
double sampleRate,
131+
int samplingMechanism,
132+
boolean rateLimiterRejected) {
133+
return setSamplingPriority(samplingPriority, rate, sampleRate, samplingMechanism);
134+
}
135+
127136
T setSpanSamplingPriority(double rate, int limit);
128137

129138
T setMetric(CharSequence name, int value);

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

Lines changed: 16 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -648,14 +648,23 @@ public final DDSpan setSamplingPriority(final int newPriority, int samplingMecha
648648
@Override
649649
public DDSpan setSamplingPriority(
650650
int samplingPriority, CharSequence rate, double sampleRate, int samplingMechanism) {
651-
if (context.setSamplingPriority(samplingPriority, samplingMechanism)) {
651+
return setSamplingPriority(samplingPriority, rate, sampleRate, samplingMechanism, false);
652+
}
653+
654+
@Override
655+
public DDSpan setSamplingPriority(
656+
int samplingPriority,
657+
CharSequence rate,
658+
double sampleRate,
659+
int samplingMechanism,
660+
boolean rateLimiterRejected) {
661+
if (context.setSamplingPriority(
662+
samplingPriority,
663+
samplingMechanism,
664+
sampleRate,
665+
getTraceId().toLong(),
666+
rateLimiterRejected)) {
652667
setMetric(rate, sampleRate);
653-
if (samplingMechanism == SamplingMechanism.AGENT_RATE
654-
|| samplingMechanism == SamplingMechanism.LOCAL_USER_RULE
655-
|| samplingMechanism == SamplingMechanism.REMOTE_USER_RULE
656-
|| samplingMechanism == SamplingMechanism.REMOTE_ADAPTIVE_RULE) {
657-
context.getPropagationTags().updateKnuthSamplingRate(sampleRate);
658-
}
659668
}
660669
return this;
661670
}

0 commit comments

Comments
 (0)