Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
@@ -0,0 +1,89 @@
import akka.actor.ActorRef
import akka.actor.ActorSystem
import akka.actor.DeadLetter
import akka.actor.Props
import akka.actor.UntypedActor
import akka.testkit.TestKit
import akka.testkit.TestProbe
import datadog.trace.agent.test.InstrumentationSpecification
import scala.concurrent.duration.Duration

import java.util.concurrent.CountDownLatch
import java.util.concurrent.atomic.AtomicInteger

import static datadog.trace.agent.test.utils.TraceUtils.basicSpan
import static datadog.trace.agent.test.utils.TraceUtils.runUnderTrace
import static java.util.concurrent.TimeUnit.SECONDS

class AkkaDeadLetterTest extends InstrumentationSpecification {
def "release context for #count queued messages discarded on actor termination (dead letter: #deadLetter)"() {
setup:
def system = ActorSystem.create("dead-letter-test")
def entered = new CountDownLatch(1)
def stop = new CountDownLatch(1)
def processed = new AtomicInteger()
def actor = system.actorOf(Props.create(StoppingActor, entered, stop, processed))
def watcher = new TestProbe(system)
watcher.watch(actor)
actor.tell("stop", ActorRef.noSender())
assert entered.await(10, SECONDS)

when:
runUnderTrace("parent") {
count.times {
def message = deadLetter ? new DeadLetter("discarded", system.deadLetters(), actor) : "discarded"
actor.tell(message, ActorRef.noSender())
}
}

then:
TEST_WRITER.empty

when:
stop.countDown()
watcher.expectTerminated(actor, Duration.create(10, SECONDS))

then:
assertTraces(1) {
trace(1) {
basicSpan(it, "parent")
}
}
processed.get() == 0

cleanup:
stop?.countDown()
TestKit.shutdownActorSystem(system, Duration.create(10, SECONDS), true)

where:
count | deadLetter
1 | false
10 | false
1 | true
}

static class StoppingActor extends UntypedActor {
private final CountDownLatch entered
private final CountDownLatch stop
private final AtomicInteger processed

StoppingActor(CountDownLatch entered, CountDownLatch stop, AtomicInteger processed) {
this.entered = entered
this.stop = stop
this.processed = processed
}

@Override
void onReceive(Object message) {
if (message == "stop") {
entered.countDown()
if (!stop.await(10, SECONDS)) {
throw new IllegalStateException("Timed out waiting to stop the actor")
}
context.stop(self)
} else {
processed.incrementAndGet()
}
}
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,70 @@
package datadog.trace.instrumentation.akka.concurrent;

import static datadog.trace.agent.tooling.bytebuddy.matcher.HierarchyMatchers.implementsInterface;
import static datadog.trace.agent.tooling.bytebuddy.matcher.NameMatchers.nameStartsWith;
import static datadog.trace.agent.tooling.bytebuddy.matcher.NameMatchers.named;
import static datadog.trace.bootstrap.instrumentation.java.concurrent.AdviceUtils.cancelTask;
import static java.util.Collections.singletonMap;
import static net.bytebuddy.matcher.ElementMatchers.isMethod;
import static net.bytebuddy.matcher.ElementMatchers.isPublic;
import static net.bytebuddy.matcher.ElementMatchers.returns;
import static net.bytebuddy.matcher.ElementMatchers.takesArgument;
import static net.bytebuddy.matcher.ElementMatchers.takesArguments;

import akka.dispatch.Envelope;
import com.google.auto.service.AutoService;
import datadog.trace.agent.tooling.Instrumenter;
import datadog.trace.agent.tooling.InstrumenterModule;
import datadog.trace.bootstrap.InstrumentationContext;
import datadog.trace.bootstrap.instrumentation.java.concurrent.State;
import java.util.Map;
import net.bytebuddy.asm.Advice;
import net.bytebuddy.description.type.TypeDescription;
import net.bytebuddy.matcher.ElementMatcher;

@AutoService(InstrumenterModule.class)
public class AkkaDeadLetterQueueInstrumentation extends InstrumenterModule.ContextTracking
implements Instrumenter.ForTypeHierarchy, Instrumenter.HasMethodAdvice {

public AkkaDeadLetterQueueInstrumentation() {
super("akka_actor_receive", "akka_actor", "akka_concurrent", "java_concurrent");
}

@Override
public String hierarchyMarkerType() {
return "akka.dispatch.MessageQueue";
}

@Override
public ElementMatcher<TypeDescription> hierarchyMatcher() {
// Match the dead-letter queue without depending on Scala's anonymous-class numbering.
return nameStartsWith("akka.dispatch.Mailboxes$")
.and(implementsInterface(named("akka.dispatch.MessageQueue")));
}

@Override
public Map<String, String> contextStore() {
return singletonMap("akka.dispatch.Envelope", State.class.getName());
}

@Override
public void methodAdvice(MethodTransformer transformer) {
transformer.applyAdvice(
isMethod()
.and(isPublic())
.and(named("enqueue"))
.and(takesArguments(2))
.and(takesArgument(0, named("akka.actor.ActorRef")))
.and(takesArgument(1, named("akka.dispatch.Envelope")))
.and(returns(void.class)),
getClass().getName() + "$EnqueueAdvice");
}

public static class EnqueueAdvice {
@Advice.OnMethodExit(onThrowable = Throwable.class, suppress = Throwable.class)
public static void afterEnqueue(@Advice.Argument(1) Envelope envelope) {
// Dead-letter publication discards the original envelope without invoking its actor.
cancelTask(InstrumentationContext.get(Envelope.class, State.class), envelope);
}
}
}
Loading