diff --git a/dd-java-agent/instrumentation/cassandra/cassandra-4.0/build.gradle b/dd-java-agent/instrumentation/cassandra/cassandra-4.0/build.gradle new file mode 100644 index 00000000000..d07fb637051 --- /dev/null +++ b/dd-java-agent/instrumentation/cassandra/cassandra-4.0/build.gradle @@ -0,0 +1,28 @@ +muzzle { + pass { + group = "com.datastax.oss" + module = "java-driver-core" + versions = "[4.0.0,)" + } +} + +apply from: "$rootDir/gradle/java.gradle" + +addTestSuiteForDir('latestDepTest', 'test') + +dependencies { + compileOnly group: 'com.datastax.oss', name: 'java-driver-core', version: '4.0.0' + + testImplementation group: 'com.datastax.oss', name: 'java-driver-core', version: '4.0.0' + testImplementation group: 'com.github.jbellis', name: 'jamm', version: '0.3.3' + testImplementation (group: 'org.testcontainers', name: 'cassandra', version: libs.versions.testcontainers.get()) { + exclude group: 'com.datastax.cassandra', module: 'cassandra-driver-core' + } + + latestDepTestImplementation group: 'com.datastax.oss', name: 'java-driver-core', version: '4.+' +} + +tasks.withType(Test).configureEach { + jvmArgs '-Dtestcontainers.ryuk.disabled=true' + environment 'TESTCONTAINERS_RYUK_DISABLED', 'true' +} diff --git a/dd-java-agent/instrumentation/cassandra/cassandra-4.0/src/main/java/datadog/trace/instrumentation/cassandra4/CassandraClientDecorator.java b/dd-java-agent/instrumentation/cassandra/cassandra-4.0/src/main/java/datadog/trace/instrumentation/cassandra4/CassandraClientDecorator.java new file mode 100644 index 00000000000..1b18bba4470 --- /dev/null +++ b/dd-java-agent/instrumentation/cassandra/cassandra-4.0/src/main/java/datadog/trace/instrumentation/cassandra4/CassandraClientDecorator.java @@ -0,0 +1,93 @@ +package datadog.trace.instrumentation.cassandra4; + +import com.datastax.oss.driver.api.core.CqlSession; +import com.datastax.oss.driver.api.core.cql.ExecutionInfo; +import com.datastax.oss.driver.api.core.cql.ResultSet; +import com.datastax.oss.driver.api.core.metadata.EndPoint; +import com.datastax.oss.driver.api.core.metadata.Node; +import datadog.trace.api.naming.SpanNaming; +import datadog.trace.api.normalize.SQLNormalizer; +import datadog.trace.bootstrap.instrumentation.api.AgentSpan; +import datadog.trace.bootstrap.instrumentation.api.InternalSpanTypes; +import datadog.trace.bootstrap.instrumentation.api.Tags; +import datadog.trace.bootstrap.instrumentation.decorator.DBTypeProcessingDatabaseClientDecorator; +import datadog.trace.bootstrap.instrumentation.jdbc.DBQueryInfo; +import java.net.InetSocketAddress; +import java.net.SocketAddress; + +public class CassandraClientDecorator extends DBTypeProcessingDatabaseClientDecorator { + private static final String DB_TYPE = "cassandra"; + private static final String SERVICE_NAME = + SpanNaming.instance().namingSchema().database().service(DB_TYPE); + public static final String OPERATION_NAME = + SpanNaming.instance().namingSchema().database().operation(DB_TYPE); + public static final String JAVA_CASSANDRA = "java-cassandra"; + + public static final CassandraClientDecorator DECORATE = new CassandraClientDecorator(); + + @Override + protected String[] instrumentationNames() { + return new String[] {"cassandra"}; + } + + @Override + protected String service() { + return SERVICE_NAME; + } + + @Override + protected CharSequence component() { + return JAVA_CASSANDRA; + } + + @Override + protected CharSequence spanType() { + return InternalSpanTypes.CASSANDRA; + } + + @Override + protected String dbType() { + return DB_TYPE; + } + + @Override + protected String dbUser(final CqlSession session) { + return null; + } + + @Override + protected String dbInstance(final CqlSession session) { + return session.getKeyspace().map(k -> k.asCql(false)).orElse(null); + } + + @Override + protected String dbHostname(final CqlSession session) { + return ContactPointsUtil.getFirstHost(session); + } + + public void onStatement(final AgentSpan span, final CharSequence statement) { + span.setResourceName(SQLNormalizer.normalize(statement.toString()).toString()); + final CharSequence operation = DBQueryInfo.extractOperation(statement); + if (operation != null) { + span.setTag(Tags.DB_OPERATION, operation); + } + } + + public void onResponse(final AgentSpan span, final ResultSet result) { + if (result != null) { + final ExecutionInfo executionInfo = result.getExecutionInfo(); + if (executionInfo != null) { + final Node coordinator = executionInfo.getCoordinator(); + if (coordinator != null) { + final EndPoint endPoint = coordinator.getEndPoint(); + if (endPoint != null) { + final SocketAddress socketAddress = endPoint.resolve(); + if (socketAddress instanceof InetSocketAddress) { + onPeerConnection(span, (InetSocketAddress) socketAddress); + } + } + } + } + } + } +} diff --git a/dd-java-agent/instrumentation/cassandra/cassandra-4.0/src/main/java/datadog/trace/instrumentation/cassandra4/CassandraClientModule.java b/dd-java-agent/instrumentation/cassandra/cassandra-4.0/src/main/java/datadog/trace/instrumentation/cassandra4/CassandraClientModule.java new file mode 100644 index 00000000000..5c0ac46ebb3 --- /dev/null +++ b/dd-java-agent/instrumentation/cassandra/cassandra-4.0/src/main/java/datadog/trace/instrumentation/cassandra4/CassandraClientModule.java @@ -0,0 +1,30 @@ +package datadog.trace.instrumentation.cassandra4; + +import com.google.auto.service.AutoService; +import datadog.trace.agent.tooling.Instrumenter; +import datadog.trace.agent.tooling.InstrumenterModule; +import java.util.Collections; +import java.util.List; + +@AutoService(InstrumenterModule.class) +public class CassandraClientModule extends InstrumenterModule.Tracing { + + public CassandraClientModule() { + super("cassandra"); + } + + @Override + public String[] helperClassNames() { + return new String[] { + packageName + ".CassandraClientDecorator", + packageName + ".CassandraDBMUtil", + packageName + ".ContactPointsUtil", + packageName + ".SpanFinishingCallback", + }; + } + + @Override + public List typeInstrumentations() { + return Collections.singletonList(new CqlSessionExecuteInstrumentation()); + } +} diff --git a/dd-java-agent/instrumentation/cassandra/cassandra-4.0/src/main/java/datadog/trace/instrumentation/cassandra4/CassandraDBMUtil.java b/dd-java-agent/instrumentation/cassandra/cassandra-4.0/src/main/java/datadog/trace/instrumentation/cassandra4/CassandraDBMUtil.java new file mode 100644 index 00000000000..736cabe67c0 --- /dev/null +++ b/dd-java-agent/instrumentation/cassandra/cassandra-4.0/src/main/java/datadog/trace/instrumentation/cassandra4/CassandraDBMUtil.java @@ -0,0 +1,65 @@ +package datadog.trace.instrumentation.cassandra4; + +import static datadog.trace.api.Config.DBM_PROPAGATION_MODE_FULL; +import static datadog.trace.bootstrap.instrumentation.api.InstrumentationTags.DBM_TRACE_INJECTED; + +import datadog.trace.api.Config; +import datadog.trace.api.propagation.W3CTraceParent; +import datadog.trace.bootstrap.instrumentation.api.AgentSpan; +import datadog.trace.bootstrap.instrumentation.dbm.SharedDBCommenter; + +/** + * Utility class for Cassandra Database Monitoring (DBM) comment injection. When DBM propagation is + * enabled, this injects trace context as a CQL comment prepended to the query string so that the + * Datadog database agent can correlate queries back to traces. + */ +public final class CassandraDBMUtil { + + private CassandraDBMUtil() {} + + /** + * Injects a DBM trace comment into the CQL query string if DBM propagation is enabled. Sets the + * {@code _dd.dbm_trace_injected} tag on the span when injection occurs. + * + * @param span the current agent span + * @param query the original CQL query string + * @param hostname the database host (may be null) + * @param dbName the database/keyspace name (may be null) + * @return the query string with DBM comment prepended, or the original query if DBM is disabled + */ + public static String injectComment(AgentSpan span, String query, String hostname, String dbName) { + if (!Config.get().isDbmCommentInjectionEnabled()) { + return query; + } + + if (query == null || query.isEmpty()) { + return query; + } + + if (span.forceSamplingDecision() == null) { + return query; + } + + String dbService = span.getServiceName(); + String traceParent = + Config.get().getDbmPropagationMode().equals(DBM_PROPAGATION_MODE_FULL) + ? W3CTraceParent.from(span) + : null; + + String commentContent = + SharedDBCommenter.buildComment(dbService, "cassandra", hostname, dbName, traceParent); + if (commentContent == null || commentContent.isEmpty()) { + return query; + } + + // Check for duplicate injection + if (SharedDBCommenter.containsTraceComment(query)) { + return query; + } + + span.setTag(DBM_TRACE_INJECTED, true); + + // Prepend the DBM comment as a CQL comment + return "/* " + commentContent + " */ " + query; + } +} diff --git a/dd-java-agent/instrumentation/cassandra/cassandra-4.0/src/main/java/datadog/trace/instrumentation/cassandra4/ContactPointsUtil.java b/dd-java-agent/instrumentation/cassandra/cassandra-4.0/src/main/java/datadog/trace/instrumentation/cassandra4/ContactPointsUtil.java new file mode 100644 index 00000000000..82b065bb818 --- /dev/null +++ b/dd-java-agent/instrumentation/cassandra/cassandra-4.0/src/main/java/datadog/trace/instrumentation/cassandra4/ContactPointsUtil.java @@ -0,0 +1,85 @@ +package datadog.trace.instrumentation.cassandra4; + +import com.datastax.oss.driver.api.core.CqlIdentifier; +import com.datastax.oss.driver.api.core.CqlSession; +import com.datastax.oss.driver.api.core.cql.BoundStatement; +import com.datastax.oss.driver.api.core.cql.Statement; +import com.datastax.oss.driver.api.core.metadata.EndPoint; +import com.datastax.oss.driver.api.core.metadata.Node; +import java.net.InetSocketAddress; +import java.net.SocketAddress; +import java.util.Collection; +import java.util.Optional; + +public final class ContactPointsUtil { + + private ContactPointsUtil() {} + + public static String getContactPoints(final CqlSession session) { + try { + Collection nodes = session.getMetadata().getNodes().values(); + if (nodes.isEmpty()) { + return null; + } + StringBuilder sb = new StringBuilder(); + for (Node node : nodes) { + if (sb.length() > 0) { + sb.append(","); + } + EndPoint endPoint = node.getEndPoint(); + if (endPoint != null) { + SocketAddress socketAddress = endPoint.resolve(); + if (socketAddress instanceof InetSocketAddress) { + InetSocketAddress inetSocketAddress = (InetSocketAddress) socketAddress; + sb.append(inetSocketAddress.getHostString()) + .append(":") + .append(inetSocketAddress.getPort()); + } + } + } + return sb.length() > 0 ? sb.toString() : null; + } catch (Throwable ignored) { + return null; + } + } + + public static String getFirstHost(final CqlSession session) { + try { + Collection nodes = session.getMetadata().getNodes().values(); + for (Node node : nodes) { + EndPoint endPoint = node.getEndPoint(); + if (endPoint != null) { + SocketAddress socketAddress = endPoint.resolve(); + if (socketAddress instanceof InetSocketAddress) { + return ((InetSocketAddress) socketAddress).getHostString(); + } + } + } + } catch (Throwable ignored) { + // Connection metadata may not be available + } + return null; + } + + public static String getKeyspace(final CqlSession session) { + try { + Optional keyspace = session.getKeyspace(); + if (keyspace.isPresent()) { + return keyspace.get().asCql(false); + } + } catch (Throwable ignored) { + // Keyspace may not be available + } + return null; + } + + public static String getQuery(final Statement statement) { + String query = null; + if (statement instanceof com.datastax.oss.driver.api.core.cql.SimpleStatement) { + query = ((com.datastax.oss.driver.api.core.cql.SimpleStatement) statement).getQuery(); + } else if (statement instanceof BoundStatement) { + query = ((BoundStatement) statement).getPreparedStatement().getQuery(); + } + return query == null ? "" : query; + } +} diff --git a/dd-java-agent/instrumentation/cassandra/cassandra-4.0/src/main/java/datadog/trace/instrumentation/cassandra4/CqlSessionExecuteInstrumentation.java b/dd-java-agent/instrumentation/cassandra/cassandra-4.0/src/main/java/datadog/trace/instrumentation/cassandra4/CqlSessionExecuteInstrumentation.java new file mode 100644 index 00000000000..595dccc8891 --- /dev/null +++ b/dd-java-agent/instrumentation/cassandra/cassandra-4.0/src/main/java/datadog/trace/instrumentation/cassandra4/CqlSessionExecuteInstrumentation.java @@ -0,0 +1,124 @@ +package datadog.trace.instrumentation.cassandra4; + +import static datadog.trace.agent.tooling.bytebuddy.matcher.NameMatchers.named; +import static datadog.trace.bootstrap.instrumentation.api.AgentTracer.activateSpan; +import static datadog.trace.bootstrap.instrumentation.api.AgentTracer.startSpan; +import static datadog.trace.instrumentation.cassandra4.CassandraClientDecorator.DECORATE; +import static datadog.trace.instrumentation.cassandra4.CassandraClientDecorator.JAVA_CASSANDRA; +import static datadog.trace.instrumentation.cassandra4.CassandraClientDecorator.OPERATION_NAME; +import static net.bytebuddy.matcher.ElementMatchers.isMethod; +import static net.bytebuddy.matcher.ElementMatchers.isPublic; +import static net.bytebuddy.matcher.ElementMatchers.takesArgument; +import static net.bytebuddy.matcher.ElementMatchers.takesArguments; + +import com.datastax.oss.driver.api.core.CqlSession; +import com.datastax.oss.driver.api.core.cql.ResultSet; +import com.datastax.oss.driver.api.core.cql.SimpleStatement; +import com.datastax.oss.driver.api.core.cql.Statement; +import com.datastax.oss.driver.api.core.session.Request; +import com.datastax.oss.driver.api.core.type.reflect.GenericType; +import datadog.trace.agent.tooling.Instrumenter; +import datadog.trace.bootstrap.CallDepthThreadLocalMap; +import datadog.trace.bootstrap.instrumentation.api.AgentScope; +import datadog.trace.bootstrap.instrumentation.api.AgentSpan; +import datadog.trace.bootstrap.instrumentation.api.InstrumentationTags; +import java.util.concurrent.CompletionStage; +import net.bytebuddy.asm.Advice; + +/** + * Instruments {@code DefaultSession.execute(Request, GenericType)} which is the single concrete + * dispatch method that all CQL operations (sync execute, async executeAsync, prepare, etc.) flow + * through. The default interface methods on {@code CqlSession} (execute, executeAsync) all delegate + * to this method, and it is concretely declared on {@code DefaultSession}, making it compatible + * with the simple method graph compiler used by default. + */ +public class CqlSessionExecuteInstrumentation + implements Instrumenter.ForSingleType, Instrumenter.HasMethodAdvice { + + @Override + public String instrumentedType() { + return "com.datastax.oss.driver.internal.core.session.DefaultSession"; + } + + @Override + public void methodAdvice(MethodTransformer transformer) { + transformer.applyAdvice( + isMethod() + .and(isPublic()) + .and(named("execute")) + .and(takesArguments(2)) + .and(takesArgument(0, named("com.datastax.oss.driver.api.core.session.Request"))) + .and( + takesArgument( + 1, named("com.datastax.oss.driver.api.core.type.reflect.GenericType"))), + CqlSessionExecuteInstrumentation.class.getName() + "$CqlSessionExecuteAdvice"); + } + + public static class CqlSessionExecuteAdvice { + + @Advice.OnMethodEnter(suppress = Throwable.class) + public static AgentScope onEnter( + @Advice.This final CqlSession session, + @Advice.Argument(value = 0, readOnly = false) Request request, + @Advice.Argument(1) final GenericType resultType) { + if (CallDepthThreadLocalMap.incrementCallDepth(CqlSession.class) > 0) { + return null; + } + // Only instrument CQL statement execution, not prepare requests or other request types + if (!(request instanceof Statement)) { + CallDepthThreadLocalMap.reset(CqlSession.class); + return null; + } + final String query = ContactPointsUtil.getQuery((Statement) request); + final AgentSpan span = startSpan(JAVA_CASSANDRA, OPERATION_NAME); + DECORATE.afterStart(span); + DECORATE.onConnection(span, session); + DECORATE.onStatement(span, query); + final String contactPoints = ContactPointsUtil.getContactPoints(session); + if (contactPoints != null) { + span.setTag(InstrumentationTags.CASSANDRA_CONTACT_POINTS, contactPoints); + } + // DBM comment injection: inject trace context into CQL for SimpleStatement only + if (request instanceof SimpleStatement) { + final String dbName = ContactPointsUtil.getKeyspace(session); + final String hostname = ContactPointsUtil.getFirstHost(session); + final String injectedQuery = CassandraDBMUtil.injectComment(span, query, hostname, dbName); + if (injectedQuery != null && !injectedQuery.equals(query)) { + request = ((SimpleStatement) request).setQuery(injectedQuery); + } + } + return activateSpan(span); + } + + @Advice.OnMethodExit(onThrowable = Throwable.class, suppress = Throwable.class) + public static void onExit( + @Advice.Enter final AgentScope scope, + @Advice.Return final Object result, + @Advice.Thrown final Throwable throwable) { + if (scope == null) { + return; + } + CallDepthThreadLocalMap.reset(CqlSession.class); + final AgentSpan span = scope.span(); + if (result instanceof CompletionStage) { + // Async path: close the scope but let the callback finish the span + scope.close(); + ((CompletionStage) result).whenComplete(new SpanFinishingCallback(span)); + } else { + // Sync path: finish span immediately + try { + if (throwable != null) { + DECORATE.onError(span, throwable); + } + if (result instanceof ResultSet) { + DECORATE.onResponse(span, (ResultSet) result); + } + DECORATE.beforeFinish(span); + } finally { + scope.close(); + span.finish(); + } + } + } + } +} diff --git a/dd-java-agent/instrumentation/cassandra/cassandra-4.0/src/main/java/datadog/trace/instrumentation/cassandra4/SpanFinishingCallback.java b/dd-java-agent/instrumentation/cassandra/cassandra-4.0/src/main/java/datadog/trace/instrumentation/cassandra4/SpanFinishingCallback.java new file mode 100644 index 00000000000..ac54b0102ca --- /dev/null +++ b/dd-java-agent/instrumentation/cassandra/cassandra-4.0/src/main/java/datadog/trace/instrumentation/cassandra4/SpanFinishingCallback.java @@ -0,0 +1,51 @@ +package datadog.trace.instrumentation.cassandra4; + +import static datadog.trace.instrumentation.cassandra4.CassandraClientDecorator.DECORATE; + +import com.datastax.oss.driver.api.core.cql.AsyncResultSet; +import com.datastax.oss.driver.api.core.cql.ExecutionInfo; +import com.datastax.oss.driver.api.core.metadata.EndPoint; +import com.datastax.oss.driver.api.core.metadata.Node; +import datadog.trace.bootstrap.instrumentation.api.AgentSpan; +import java.net.InetSocketAddress; +import java.net.SocketAddress; +import java.util.function.BiConsumer; + +public class SpanFinishingCallback implements BiConsumer { + + private final AgentSpan span; + + public SpanFinishingCallback(final AgentSpan span) { + this.span = span; + } + + @Override + public void accept(final Object result, final Throwable error) { + try { + if (error != null) { + DECORATE.onError(span, error); + } else if (result instanceof AsyncResultSet) { + onAsyncResponse(span, (AsyncResultSet) result); + } + DECORATE.beforeFinish(span); + } finally { + span.finish(); + } + } + + private static void onAsyncResponse(final AgentSpan span, final AsyncResultSet result) { + final ExecutionInfo executionInfo = result.getExecutionInfo(); + if (executionInfo != null) { + final Node coordinator = executionInfo.getCoordinator(); + if (coordinator != null) { + final EndPoint endPoint = coordinator.getEndPoint(); + if (endPoint != null) { + final SocketAddress socketAddress = endPoint.resolve(); + if (socketAddress instanceof InetSocketAddress) { + DECORATE.onPeerConnection(span, (InetSocketAddress) socketAddress); + } + } + } + } + } +} diff --git a/dd-java-agent/instrumentation/cassandra/cassandra-4.0/src/test/java/datadog/trace/instrumentation/cassandra4/CassandraClientTest.java b/dd-java-agent/instrumentation/cassandra/cassandra-4.0/src/test/java/datadog/trace/instrumentation/cassandra4/CassandraClientTest.java new file mode 100644 index 00000000000..7b36332d6d9 --- /dev/null +++ b/dd-java-agent/instrumentation/cassandra/cassandra-4.0/src/test/java/datadog/trace/instrumentation/cassandra4/CassandraClientTest.java @@ -0,0 +1,475 @@ +package datadog.trace.instrumentation.cassandra4; + +import static datadog.trace.agent.test.assertions.TagsMatcher.defaultTags; +import static datadog.trace.agent.test.assertions.TagsMatcher.error; +import static datadog.trace.agent.test.assertions.TagsMatcher.tag; +import static datadog.trace.agent.test.assertions.TraceMatcher.trace; +import static datadog.trace.api.config.TraceInstrumentationConfig.DB_DBM_PROPAGATION_MODE_MODE; +import static datadog.trace.test.junit.utils.assertions.Matchers.any; +import static datadog.trace.test.junit.utils.assertions.Matchers.is; +import static datadog.trace.test.junit.utils.assertions.Matchers.isNonNull; +import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import com.datastax.oss.driver.api.core.CqlSession; +import com.datastax.oss.driver.api.core.cql.AsyncResultSet; +import com.datastax.oss.driver.api.core.cql.BoundStatement; +import com.datastax.oss.driver.api.core.cql.PreparedStatement; +import com.datastax.oss.driver.api.core.cql.ResultSet; +import com.datastax.oss.driver.api.core.servererrors.SyntaxError; +import datadog.trace.agent.test.AbstractInstrumentationTest; +import datadog.trace.agent.test.assertions.SpanMatcher; +import datadog.trace.api.DDSpanTypes; +import datadog.trace.bootstrap.instrumentation.api.InstrumentationTags; +import datadog.trace.bootstrap.instrumentation.api.Tags; +import datadog.trace.test.junit.utils.config.WithConfig; +import java.net.InetSocketAddress; +import java.time.Duration; +import java.util.concurrent.CompletionStage; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; +import org.junit.jupiter.api.AfterAll; +import org.junit.jupiter.api.BeforeAll; +import org.junit.jupiter.api.MethodOrderer; +import org.junit.jupiter.api.Order; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.TestMethodOrder; +import org.testcontainers.containers.CassandraContainer; + +@TestMethodOrder(MethodOrderer.OrderAnnotation.class) +class CassandraClientTest extends AbstractInstrumentationTest { + + private static final String OPERATION = "cassandra.query"; + private static final String SERVICE = "cassandra"; + private static final String COMPONENT = "java-cassandra"; + + private static CassandraContainer container; + private static CqlSession session; + private static int port; + private static String contactPoint; + + @BeforeAll + static void setupCassandra() { + container = new CassandraContainer<>("cassandra:3").withStartupTimeout(Duration.ofSeconds(120)); + container.start(); + port = container.getMappedPort(9042); + InetSocketAddress contact = container.getContactPoint(); + contactPoint = contact.getHostString() + ":" + contact.getPort(); + session = + CqlSession.builder() + .addContactPoint(contact) + .withLocalDatacenter(container.getLocalDatacenter()) + .build(); + + // Wait for connection setup traces and clear them + try { + writer.waitForTraces(1); + } catch (Exception ignored) { + } + tracer.flush(); + writer.start(); + } + + @AfterAll + static void tearDownCassandra() { + if (session != null) { + session.close(); + } + if (container != null) { + container.stop(); + } + } + + @Test + @Order(1) + void syncCreateKeyspace() { + session.execute( + "CREATE KEYSPACE IF NOT EXISTS sync_test WITH REPLICATION = {'class':'SimpleStrategy', 'replication_factor':3}"); + + assertTraces( + trace( + cassandraSpan( + "CREATE KEYSPACE IF NOT EXISTS sync_test WITH REPLICATION = {'class':'SimpleStrategy', 'replication_factor':3}", + null))); + } + + @Test + @Order(2) + void syncCreateTable() { + session.execute("CREATE TABLE IF NOT EXISTS sync_test.users (id UUID PRIMARY KEY, name text)"); + + assertTraces( + trace( + cassandraSpan( + "CREATE TABLE IF NOT EXISTS sync_test.users (id UUID PRIMARY KEY, name text)", + null))); + } + + @Test + @Order(3) + void syncInsert() { + session.execute("INSERT INTO sync_test.users (id, name) values (uuid(), 'alice')"); + + assertTraces( + trace(cassandraSpan("INSERT INTO sync_test.users (id, name) values (uuid(), ?)", null))); + } + + @Test + @Order(4) + void syncSelect() { + ResultSet result = + session.execute("SELECT * FROM sync_test.users where name = 'alice' ALLOW FILTERING"); + assertNotNull(result); + + assertTraces( + trace(cassandraSpan("SELECT * FROM sync_test.users where name = ? ALLOW FILTERING", null))); + } + + @Test + @Order(5) + void asyncQuery() throws Exception { + CountDownLatch latch = new CountDownLatch(1); + CompletionStage future = + session.executeAsync("SELECT * FROM sync_test.users where name = 'alice' ALLOW FILTERING"); + future.whenComplete( + (result, error) -> { + latch.countDown(); + }); + + assertTrue(latch.await(10, TimeUnit.SECONDS), "Async query did not complete in time"); + + assertTraces( + trace(cassandraSpan("SELECT * FROM sync_test.users where name = ? ALLOW FILTERING", null))); + } + + @Test + @Order(6) + void syncQueryError() { + assertThrows( + SyntaxError.class, + () -> { + session.execute("INVALID CQL QUERY"); + }); + + assertTraces(trace(cassandraSpanWithError("INVALID CQL QUERY", null, SyntaxError.class))); + } + + @Test + @Order(7) + void asyncQueryError() throws Exception { + CountDownLatch latch = new CountDownLatch(1); + CompletionStage future = session.executeAsync("INVALID ASYNC CQL"); + future.whenComplete( + (result, error) -> { + latch.countDown(); + }); + + assertTrue(latch.await(10, TimeUnit.SECONDS), "Async query did not complete in time"); + + assertTraces(trace(cassandraSpanWithError("INVALID ASYNC CQL", null, SyntaxError.class))); + } + + @Test + @Order(8) + void syncBoundStatementExecute() { + // Prepare does not create a span (not a Statement), only execute does + PreparedStatement prepared = + session.prepare("SELECT * FROM sync_test.users where name = ? ALLOW FILTERING"); + BoundStatement bound = prepared.bind("alice"); + ResultSet result = session.execute(bound); + assertNotNull(result); + + // BoundStatement execution triggers a span with the normalized prepared query + assertTraces( + trace(cassandraSpan("SELECT * FROM sync_test.users where name = ? ALLOW FILTERING", null))); + } + + @Test + @Order(9) + void asyncBoundStatementExecute() throws Exception { + // Prepare does not create a span (not a Statement), only execute does + PreparedStatement prepared = + session.prepare("INSERT INTO sync_test.users (id, name) values (uuid(), ?)"); + BoundStatement bound = prepared.bind("bob"); + CountDownLatch latch = new CountDownLatch(1); + CompletionStage future = session.executeAsync(bound); + future.whenComplete( + (result, error) -> { + latch.countDown(); + }); + + assertTrue(latch.await(10, TimeUnit.SECONDS), "Async prepared query did not complete in time"); + + // Async BoundStatement execution triggers a span with the normalized prepared query + assertTraces( + trace(cassandraSpan("INSERT INTO sync_test.users (id, name) values (uuid(), ?)", null))); + } + + @Test + @Order(10) + void syncDropKeyspace() { + session.execute("DROP KEYSPACE IF EXISTS sync_test"); + + assertTraces(trace(cassandraSpan("DROP KEYSPACE IF EXISTS sync_test", null))); + } + + @Test + @Order(11) + @WithConfig(key = DB_DBM_PROPAGATION_MODE_MODE, value = "service") + void dbmServiceModeInjectsSyncQuery() { + session.execute( + "CREATE KEYSPACE IF NOT EXISTS dbm_test WITH REPLICATION = {'class':'SimpleStrategy', 'replication_factor':3}"); + + assertTraces( + trace( + cassandraSpanWithDbm( + "CREATE KEYSPACE IF NOT EXISTS dbm_test WITH REPLICATION = {'class':'SimpleStrategy', 'replication_factor':3}", + null))); + } + + @Test + @Order(12) + @WithConfig(key = DB_DBM_PROPAGATION_MODE_MODE, value = "full") + void dbmFullModeInjectsSyncQuery() { + session.execute("CREATE TABLE IF NOT EXISTS dbm_test.items (id UUID PRIMARY KEY, name text)"); + + assertTraces( + trace( + cassandraSpanWithDbm( + "CREATE TABLE IF NOT EXISTS dbm_test.items (id UUID PRIMARY KEY, name text)", + null))); + } + + @Test + @Order(13) + @WithConfig(key = DB_DBM_PROPAGATION_MODE_MODE, value = "full") + void dbmFullModeInjectsAsyncQuery() throws Exception { + CountDownLatch latch = new CountDownLatch(1); + CompletionStage future = + session.executeAsync("INSERT INTO dbm_test.items (id, name) values (uuid(), 'test_item')"); + future.whenComplete( + (result, error) -> { + latch.countDown(); + }); + + assertTrue(latch.await(10, TimeUnit.SECONDS), "Async DBM query did not complete in time"); + + assertTraces( + trace( + cassandraSpanWithDbm( + "INSERT INTO dbm_test.items (id, name) values (uuid(), ?)", null))); + } + + @Test + @Order(14) + void dbmDisabledDoesNotInject() { + // With no DBM config, _dd.dbm_trace_injected tag should not be present + session.execute("SELECT * FROM dbm_test.items LIMIT 1"); + + assertTraces(trace(cassandraSpan("SELECT * FROM dbm_test.items LIMIT ?", null))); + } + + @Test + @Order(15) + void dbmCleanup() { + session.execute("DROP KEYSPACE IF EXISTS dbm_test"); + + assertTraces(trace(cassandraSpan("DROP KEYSPACE IF EXISTS dbm_test", null))); + } + + @Test + @Order(16) + void peerServiceInputTagsSetWithoutKeyspace() { + // Without a keyspace, peer.hostname should still be set (from dbHostname()), + // but db.instance should not be present. + session.execute( + "CREATE KEYSPACE IF NOT EXISTS peer_test WITH REPLICATION = {'class':'SimpleStrategy', 'replication_factor':3}"); + + assertTraces( + trace( + cassandraSpan( + "CREATE KEYSPACE IF NOT EXISTS peer_test WITH REPLICATION = {'class':'SimpleStrategy', 'replication_factor':3}", + null))); + } + + @Test + @Order(17) + void peerServiceInputTagsSetWithKeyspace() { + // Create a session with a keyspace so db.instance is populated, + // verifying both peer.hostname and db.instance are set as peer service inputs. + InetSocketAddress contact = container.getContactPoint(); + try (CqlSession keyspaceSession = + CqlSession.builder() + .addContactPoint(contact) + .withLocalDatacenter(container.getLocalDatacenter()) + .withKeyspace("peer_test") + .build()) { + // Clear any setup traces + try { + writer.waitForTraces(1); + } catch (Exception ignored) { + } + tracer.flush(); + writer.start(); + + keyspaceSession.execute("CREATE TABLE IF NOT EXISTS users (id UUID PRIMARY KEY, name text)"); + + assertTraces( + trace( + cassandraSpanWithKeyspace( + "CREATE TABLE IF NOT EXISTS users (id UUID PRIMARY KEY, name text)", + "\"peer_test\""))); + } + } + + @Test + @Order(18) + void peerServiceCleanup() { + session.execute("DROP KEYSPACE IF EXISTS peer_test"); + + assertTraces(trace(cassandraSpan("DROP KEYSPACE IF EXISTS peer_test", null))); + } + + private static SpanMatcher cassandraSpan(String resource, String keyspace) { + SpanMatcher matcher = + SpanMatcher.span() + .root() + .serviceName(SERVICE) + .operationName(OPERATION) + .resourceName(resource) + .type(DDSpanTypes.CASSANDRA) + .measured(); + if (keyspace != null) { + matcher.tags( + tag(Tags.COMPONENT, is(COMPONENT)), + tag(Tags.SPAN_KIND, is(Tags.SPAN_KIND_CLIENT)), + tag(Tags.PEER_HOSTNAME, isNonNull()), + tag(Tags.PEER_PORT, is(port)), + tag(Tags.DB_TYPE, is("cassandra")), + tag(Tags.DB_INSTANCE, is(keyspace)), + tag(Tags.DB_OPERATION, is(resource.split(" ")[0])), + tag(InstrumentationTags.CASSANDRA_CONTACT_POINTS, is(contactPoint)), + tag("peer.ipv4", any()), + tag("_dd.svc_src", any()), + defaultTags()); + } else { + matcher.tags( + tag(Tags.COMPONENT, is(COMPONENT)), + tag(Tags.SPAN_KIND, is(Tags.SPAN_KIND_CLIENT)), + tag(Tags.PEER_HOSTNAME, isNonNull()), + tag(Tags.PEER_PORT, is(port)), + tag(Tags.DB_TYPE, is("cassandra")), + tag(Tags.DB_OPERATION, is(resource.split(" ")[0])), + tag(InstrumentationTags.CASSANDRA_CONTACT_POINTS, is(contactPoint)), + tag("peer.ipv4", any()), + tag("_dd.svc_src", any()), + defaultTags()); + } + return matcher; + } + + private static SpanMatcher cassandraSpanWithError( + String resource, String keyspace, Class errorType) { + SpanMatcher matcher = + SpanMatcher.span() + .root() + .serviceName(SERVICE) + .operationName(OPERATION) + .resourceName(resource) + .type(DDSpanTypes.CASSANDRA) + .error() + .measured(); + if (keyspace != null) { + matcher.tags( + tag(Tags.COMPONENT, is(COMPONENT)), + tag(Tags.SPAN_KIND, is(Tags.SPAN_KIND_CLIENT)), + tag(Tags.PEER_HOSTNAME, isNonNull()), + tag(Tags.DB_TYPE, is("cassandra")), + tag(Tags.DB_INSTANCE, is(keyspace)), + tag(Tags.DB_OPERATION, is(resource.split(" ")[0])), + tag(InstrumentationTags.CASSANDRA_CONTACT_POINTS, is(contactPoint)), + tag("_dd.svc_src", any()), + tag("error.message", any()), + error(errorType), + defaultTags()); + } else { + matcher.tags( + tag(Tags.COMPONENT, is(COMPONENT)), + tag(Tags.SPAN_KIND, is(Tags.SPAN_KIND_CLIENT)), + tag(Tags.PEER_HOSTNAME, isNonNull()), + tag(Tags.DB_TYPE, is("cassandra")), + tag(Tags.DB_OPERATION, is(resource.split(" ")[0])), + tag(InstrumentationTags.CASSANDRA_CONTACT_POINTS, is(contactPoint)), + tag("_dd.svc_src", any()), + tag("error.message", any()), + error(errorType), + defaultTags()); + } + return matcher; + } + + private static SpanMatcher cassandraSpanWithKeyspace(String resource, String keyspace) { + SpanMatcher matcher = + SpanMatcher.span() + .root() + .serviceName(SERVICE) + .operationName(OPERATION) + .resourceName(resource) + .type(DDSpanTypes.CASSANDRA) + .measured(); + matcher.tags( + tag(Tags.COMPONENT, is(COMPONENT)), + tag(Tags.SPAN_KIND, is(Tags.SPAN_KIND_CLIENT)), + tag(Tags.PEER_HOSTNAME, isNonNull()), + tag(Tags.PEER_PORT, is(port)), + tag(Tags.DB_TYPE, is("cassandra")), + tag(Tags.DB_INSTANCE, is(keyspace)), + tag(Tags.DB_OPERATION, is(resource.split(" ")[0])), + tag(InstrumentationTags.CASSANDRA_CONTACT_POINTS, isNonNull()), + tag("peer.ipv4", any()), + tag("_dd.svc_src", any()), + defaultTags()); + return matcher; + } + + private static SpanMatcher cassandraSpanWithDbm(String resource, String keyspace) { + SpanMatcher matcher = + SpanMatcher.span() + .root() + .serviceName(SERVICE) + .operationName(OPERATION) + .resourceName(resource) + .type(DDSpanTypes.CASSANDRA) + .measured(); + if (keyspace != null) { + matcher.tags( + tag(Tags.COMPONENT, is(COMPONENT)), + tag(Tags.SPAN_KIND, is(Tags.SPAN_KIND_CLIENT)), + tag(Tags.PEER_HOSTNAME, isNonNull()), + tag(Tags.PEER_PORT, is(port)), + tag(Tags.DB_TYPE, is("cassandra")), + tag(Tags.DB_INSTANCE, is(keyspace)), + tag(Tags.DB_OPERATION, is(resource.split(" ")[0])), + tag(InstrumentationTags.CASSANDRA_CONTACT_POINTS, is(contactPoint)), + tag(InstrumentationTags.DBM_TRACE_INJECTED, is(true)), + tag("peer.ipv4", any()), + tag("_dd.svc_src", any()), + defaultTags()); + } else { + matcher.tags( + tag(Tags.COMPONENT, is(COMPONENT)), + tag(Tags.SPAN_KIND, is(Tags.SPAN_KIND_CLIENT)), + tag(Tags.PEER_HOSTNAME, isNonNull()), + tag(Tags.PEER_PORT, is(port)), + tag(Tags.DB_TYPE, is("cassandra")), + tag(Tags.DB_OPERATION, is(resource.split(" ")[0])), + tag(InstrumentationTags.CASSANDRA_CONTACT_POINTS, is(contactPoint)), + tag(InstrumentationTags.DBM_TRACE_INJECTED, is(true)), + tag("peer.ipv4", any()), + tag("_dd.svc_src", any()), + defaultTags()); + } + return matcher; + } +} diff --git a/settings.gradle.kts b/settings.gradle.kts index 86aa23cc9d4..791e356d7df 100644 --- a/settings.gradle.kts +++ b/settings.gradle.kts @@ -321,6 +321,7 @@ include( ":dd-java-agent:instrumentation:axway-api-7.5", ":dd-java-agent:instrumentation:azure-functions-1.2.2", ":dd-java-agent:instrumentation:caffeine-1.0", + ":dd-java-agent:instrumentation:cassandra:cassandra-4.0", ":dd-java-agent:instrumentation:cdi-1.2", ":dd-java-agent:instrumentation:cics-9.1", ":dd-java-agent:instrumentation:commons-codec-1.1",