diff --git a/dd-java-agent/instrumentation/postgresql/postgresql-42.0/build.gradle b/dd-java-agent/instrumentation/postgresql/postgresql-42.0/build.gradle new file mode 100644 index 00000000000..b17b36cfb9f --- /dev/null +++ b/dd-java-agent/instrumentation/postgresql/postgresql-42.0/build.gradle @@ -0,0 +1,26 @@ +muzzle { + pass { + group = "org.postgresql" + module = "postgresql" + versions = "[42.0.0,)" + } +} + +apply from: "$rootDir/gradle/java.gradle" + +addTestSuiteForDir('latestDepTest', 'test') + +dependencies { + compileOnly group: 'org.postgresql', name: 'postgresql', version: '42.0.0' + + testImplementation group: 'org.postgresql', name: 'postgresql', version: '42.0.0' + testImplementation group: 'org.testcontainers', name: 'postgresql', version: libs.versions.testcontainers.get() + + latestDepTestImplementation group: 'org.postgresql', name: 'postgresql', version: '42.+' +} + +tasks.withType(Test).configureEach { + usesService(testcontainersLimit) + environment 'TESTCONTAINERS_RYUK_DISABLED', 'true' + jvmArgs '-Dtestcontainers.reuse.enable=false' +} diff --git a/dd-java-agent/instrumentation/postgresql/postgresql-42.0/src/main/java/datadog/trace/instrumentation/postgresql/PgPreparedStatementAdvice.java b/dd-java-agent/instrumentation/postgresql/postgresql-42.0/src/main/java/datadog/trace/instrumentation/postgresql/PgPreparedStatementAdvice.java new file mode 100644 index 00000000000..7df9e654e0e --- /dev/null +++ b/dd-java-agent/instrumentation/postgresql/postgresql-42.0/src/main/java/datadog/trace/instrumentation/postgresql/PgPreparedStatementAdvice.java @@ -0,0 +1,97 @@ +package datadog.trace.instrumentation.postgresql; + +import static datadog.trace.bootstrap.instrumentation.api.AgentTracer.activateSpan; +import static datadog.trace.bootstrap.instrumentation.api.AgentTracer.startSpan; +import static datadog.trace.instrumentation.postgresql.PostgreSQLDecorator.DECORATE; +import static datadog.trace.instrumentation.postgresql.PostgreSQLDecorator.JAVA_POSTGRESQL; +import static datadog.trace.instrumentation.postgresql.PostgreSQLDecorator.POSTGRESQL_QUERY; + +import datadog.trace.bootstrap.CallDepthThreadLocalMap; +import datadog.trace.bootstrap.InstrumentationContext; +import datadog.trace.bootstrap.instrumentation.api.AgentScope; +import datadog.trace.bootstrap.instrumentation.api.AgentSpan; +import datadog.trace.bootstrap.instrumentation.api.Tags; +import datadog.trace.bootstrap.instrumentation.jdbc.DBInfo; +import datadog.trace.bootstrap.instrumentation.jdbc.DBQueryInfo; +import java.sql.Statement; +import java.util.Collections; +import java.util.Map; +import java.util.WeakHashMap; +import net.bytebuddy.asm.Advice; + +public class PgPreparedStatementAdvice { + + public static final Map PREPARED_SQL = + Collections.synchronizedMap(new WeakHashMap()); + + public static final class ConstructorAdvice { + + @Advice.OnMethodExit(suppress = Throwable.class) + public static void onExit( + @Advice.This final Statement statement, @Advice.Argument(1) final Object query) { + // In PostgreSQL JDBC 42.0+, argument[1] is CachedQuery, use toString() to get SQL + if (query != null) { + PREPARED_SQL.put(statement, query.toString()); + } + } + } + + public static final class ExecuteAdvice { + + @Advice.OnMethodEnter(suppress = Throwable.class) + public static AgentScope onEnter(@Advice.This final Statement statement) { + final int callDepth = CallDepthThreadLocalMap.incrementCallDepth(Statement.class); + if (callDepth > 0) { + return null; + } + + final AgentSpan span = startSpan(JAVA_POSTGRESQL.toString(), POSTGRESQL_QUERY); + DECORATE.afterStart(span); + + DBInfo dbInfo = InstrumentationContext.get(Statement.class, DBInfo.class).get(statement); + if (dbInfo == null) { + dbInfo = PgStatementAdvice.ExecuteQueryAdvice.extractDbInfo(statement); + if (dbInfo != null) { + InstrumentationContext.get(Statement.class, DBInfo.class).put(statement, dbInfo); + } + } + if (dbInfo != null) { + DECORATE.onConnection(span, dbInfo); + if (dbInfo.getPort() != null) { + DECORATE.setPeerPort(span, dbInfo.getPort()); + } + } + + final String sql = PREPARED_SQL.get(statement); + if (sql != null) { + final DBQueryInfo dbQueryInfo = DBQueryInfo.ofPreparedStatement(sql); + DECORATE.onStatement(span, dbQueryInfo.getSql()); + span.setTag(Tags.DB_OPERATION, dbQueryInfo.getOperation()); + + // DBM: inject SQL comment with trace context for prepared statements + // For prepared statements, we inject via modifying the stored SQL that will + // be sent to the server. The comment is only added to spans' context, + // the actual SQL modification for prepared statements happens at connection level. + PostgreSQLSQLCommenter.inject(sql, span, dbInfo); + } + + return activateSpan(span); + } + + @Advice.OnMethodExit(onThrowable = Throwable.class, suppress = Throwable.class) + public static void onExit( + @Advice.Enter final AgentScope scope, @Advice.Thrown final Throwable throwable) { + if (scope == null) { + return; + } + final AgentSpan span = scope.span(); + if (throwable != null) { + DECORATE.onError(span, throwable); + } + DECORATE.beforeFinish(span); + scope.close(); + span.finish(); + CallDepthThreadLocalMap.reset(Statement.class); + } + } +} diff --git a/dd-java-agent/instrumentation/postgresql/postgresql-42.0/src/main/java/datadog/trace/instrumentation/postgresql/PgPreparedStatementInstrumentation.java b/dd-java-agent/instrumentation/postgresql/postgresql-42.0/src/main/java/datadog/trace/instrumentation/postgresql/PgPreparedStatementInstrumentation.java new file mode 100644 index 00000000000..98b0686123b --- /dev/null +++ b/dd-java-agent/instrumentation/postgresql/postgresql-42.0/src/main/java/datadog/trace/instrumentation/postgresql/PgPreparedStatementInstrumentation.java @@ -0,0 +1,37 @@ +package datadog.trace.instrumentation.postgresql; + +import static datadog.trace.agent.tooling.bytebuddy.matcher.NameMatchers.named; +import static net.bytebuddy.matcher.ElementMatchers.isConstructor; +import static net.bytebuddy.matcher.ElementMatchers.isMethod; +import static net.bytebuddy.matcher.ElementMatchers.isPublic; +import static net.bytebuddy.matcher.ElementMatchers.takesArguments; + +import datadog.trace.agent.tooling.Instrumenter; + +public final class PgPreparedStatementInstrumentation + implements Instrumenter.ForSingleType, Instrumenter.HasMethodAdvice { + + @Override + public String instrumentedType() { + return "org.postgresql.jdbc.PgPreparedStatement"; + } + + @Override + public void methodAdvice(MethodTransformer transformer) { + transformer.applyAdvice( + isConstructor(), + "datadog.trace.instrumentation.postgresql.PgPreparedStatementAdvice$ConstructorAdvice"); + transformer.applyAdvice( + isMethod().and(isPublic()).and(named("executeQuery")).and(takesArguments(0)), + "datadog.trace.instrumentation.postgresql.PgPreparedStatementAdvice$ExecuteAdvice"); + transformer.applyAdvice( + isMethod().and(isPublic()).and(named("executeUpdate")).and(takesArguments(0)), + "datadog.trace.instrumentation.postgresql.PgPreparedStatementAdvice$ExecuteAdvice"); + transformer.applyAdvice( + isMethod().and(isPublic()).and(named("execute")).and(takesArguments(0)), + "datadog.trace.instrumentation.postgresql.PgPreparedStatementAdvice$ExecuteAdvice"); + transformer.applyAdvice( + isMethod().and(isPublic()).and(named("executeBatch")).and(takesArguments(0)), + "datadog.trace.instrumentation.postgresql.PgPreparedStatementAdvice$ExecuteAdvice"); + } +} diff --git a/dd-java-agent/instrumentation/postgresql/postgresql-42.0/src/main/java/datadog/trace/instrumentation/postgresql/PgStatementAdvice.java b/dd-java-agent/instrumentation/postgresql/postgresql-42.0/src/main/java/datadog/trace/instrumentation/postgresql/PgStatementAdvice.java new file mode 100644 index 00000000000..33ee76e57e8 --- /dev/null +++ b/dd-java-agent/instrumentation/postgresql/postgresql-42.0/src/main/java/datadog/trace/instrumentation/postgresql/PgStatementAdvice.java @@ -0,0 +1,165 @@ +package datadog.trace.instrumentation.postgresql; + +import static datadog.trace.bootstrap.instrumentation.api.AgentTracer.activateSpan; +import static datadog.trace.bootstrap.instrumentation.api.AgentTracer.startSpan; +import static datadog.trace.instrumentation.postgresql.PostgreSQLDecorator.DECORATE; +import static datadog.trace.instrumentation.postgresql.PostgreSQLDecorator.JAVA_POSTGRESQL; +import static datadog.trace.instrumentation.postgresql.PostgreSQLDecorator.POSTGRESQL_QUERY; + +import datadog.trace.bootstrap.CallDepthThreadLocalMap; +import datadog.trace.bootstrap.InstrumentationContext; +import datadog.trace.bootstrap.instrumentation.api.AgentScope; +import datadog.trace.bootstrap.instrumentation.api.AgentSpan; +import datadog.trace.bootstrap.instrumentation.api.Tags; +import datadog.trace.bootstrap.instrumentation.jdbc.DBInfo; +import datadog.trace.bootstrap.instrumentation.jdbc.DBQueryInfo; +import datadog.trace.bootstrap.instrumentation.jdbc.JDBCConnectionUrlParser; +import java.sql.Connection; +import java.sql.DatabaseMetaData; +import java.sql.Statement; +import java.util.Collections; +import java.util.Map; +import java.util.WeakHashMap; +import net.bytebuddy.asm.Advice; + +public class PgStatementAdvice { + + public static final Map BATCH_SQL = + Collections.synchronizedMap(new WeakHashMap()); + + public static final class ExecuteQueryAdvice { + + @Advice.OnMethodEnter(suppress = Throwable.class) + public static AgentScope onEnter( + @Advice.This final Statement statement, + @Advice.Argument(value = 0, readOnly = false) String sql) { + final int callDepth = CallDepthThreadLocalMap.incrementCallDepth(Statement.class); + if (callDepth > 0) { + return null; + } + + final AgentSpan span = startSpan(JAVA_POSTGRESQL.toString(), POSTGRESQL_QUERY); + DECORATE.afterStart(span); + + DBInfo dbInfo = InstrumentationContext.get(Statement.class, DBInfo.class).get(statement); + if (dbInfo == null) { + dbInfo = extractDbInfo(statement); + if (dbInfo != null) { + InstrumentationContext.get(Statement.class, DBInfo.class).put(statement, dbInfo); + } + } + if (dbInfo != null) { + DECORATE.onConnection(span, dbInfo); + if (dbInfo.getPort() != null) { + DECORATE.setPeerPort(span, dbInfo.getPort()); + } + } + + final DBQueryInfo dbQueryInfo = DBQueryInfo.ofStatement(sql); + DECORATE.onStatement(span, dbQueryInfo.getSql()); + span.setTag(Tags.DB_OPERATION, dbQueryInfo.getOperation()); + + // DBM: inject SQL comment with trace context + sql = PostgreSQLSQLCommenter.inject(sql, span, dbInfo); + + return activateSpan(span); + } + + @Advice.OnMethodExit(onThrowable = Throwable.class, suppress = Throwable.class) + public static void onExit( + @Advice.Enter final AgentScope scope, @Advice.Thrown final Throwable throwable) { + if (scope == null) { + return; + } + final AgentSpan span = scope.span(); + if (throwable != null) { + DECORATE.onError(span, throwable); + } + DECORATE.beforeFinish(span); + scope.close(); + span.finish(); + CallDepthThreadLocalMap.reset(Statement.class); + } + + public static DBInfo extractDbInfo(final Statement statement) { + try { + final Connection connection = statement.getConnection(); + if (connection != null) { + final DatabaseMetaData metaData = connection.getMetaData(); + final String url = metaData.getURL(); + final String user = metaData.getUserName(); + DBInfo dbInfo = JDBCConnectionUrlParser.parse(url, null); + if (user != null && !user.isEmpty()) { + dbInfo = dbInfo.toBuilder().user(user).build(); + } + return dbInfo; + } + } catch (final Exception ignored) { + // Unable to extract connection info + } + return null; + } + } + + public static final class AddBatchAdvice { + + @Advice.OnMethodEnter(suppress = Throwable.class) + public static void onEnter( + @Advice.This final Statement statement, @Advice.Argument(0) final String sql) { + BATCH_SQL.put(statement, sql); + } + } + + public static final class ExecuteBatchAdvice { + + @Advice.OnMethodEnter(suppress = Throwable.class) + public static AgentScope onEnter(@Advice.This final Statement statement) { + final int callDepth = CallDepthThreadLocalMap.incrementCallDepth(Statement.class); + if (callDepth > 0) { + return null; + } + + final AgentSpan span = startSpan(JAVA_POSTGRESQL.toString(), POSTGRESQL_QUERY); + DECORATE.afterStart(span); + + DBInfo dbInfo = InstrumentationContext.get(Statement.class, DBInfo.class).get(statement); + if (dbInfo == null) { + dbInfo = ExecuteQueryAdvice.extractDbInfo(statement); + if (dbInfo != null) { + InstrumentationContext.get(Statement.class, DBInfo.class).put(statement, dbInfo); + } + } + if (dbInfo != null) { + DECORATE.onConnection(span, dbInfo); + if (dbInfo.getPort() != null) { + DECORATE.setPeerPort(span, dbInfo.getPort()); + } + } + + final String batchSql = BATCH_SQL.get(statement); + if (batchSql != null) { + final DBQueryInfo dbQueryInfo = DBQueryInfo.ofStatement(batchSql); + DECORATE.onStatement(span, dbQueryInfo.getSql()); + span.setTag(Tags.DB_OPERATION, dbQueryInfo.getOperation()); + } + + return activateSpan(span); + } + + @Advice.OnMethodExit(onThrowable = Throwable.class, suppress = Throwable.class) + public static void onExit( + @Advice.Enter final AgentScope scope, @Advice.Thrown final Throwable throwable) { + if (scope == null) { + return; + } + final AgentSpan span = scope.span(); + if (throwable != null) { + DECORATE.onError(span, throwable); + } + DECORATE.beforeFinish(span); + scope.close(); + span.finish(); + CallDepthThreadLocalMap.reset(Statement.class); + } + } +} diff --git a/dd-java-agent/instrumentation/postgresql/postgresql-42.0/src/main/java/datadog/trace/instrumentation/postgresql/PgStatementInstrumentation.java b/dd-java-agent/instrumentation/postgresql/postgresql-42.0/src/main/java/datadog/trace/instrumentation/postgresql/PgStatementInstrumentation.java new file mode 100644 index 00000000000..a05f94b9f3d --- /dev/null +++ b/dd-java-agent/instrumentation/postgresql/postgresql-42.0/src/main/java/datadog/trace/instrumentation/postgresql/PgStatementInstrumentation.java @@ -0,0 +1,53 @@ +package datadog.trace.instrumentation.postgresql; + +import static datadog.trace.agent.tooling.bytebuddy.matcher.NameMatchers.named; +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 datadog.trace.agent.tooling.Instrumenter; + +public final class PgStatementInstrumentation + implements Instrumenter.ForSingleType, Instrumenter.HasMethodAdvice { + + @Override + public String instrumentedType() { + return "org.postgresql.jdbc.PgStatement"; + } + + @Override + public void methodAdvice(MethodTransformer transformer) { + transformer.applyAdvice( + isMethod() + .and(isPublic()) + .and(named("executeQuery")) + .and(takesArguments(1)) + .and(takesArgument(0, String.class)), + "datadog.trace.instrumentation.postgresql.PgStatementAdvice$ExecuteQueryAdvice"); + transformer.applyAdvice( + isMethod() + .and(isPublic()) + .and(named("executeUpdate")) + .and(takesArguments(1)) + .and(takesArgument(0, String.class)), + "datadog.trace.instrumentation.postgresql.PgStatementAdvice$ExecuteQueryAdvice"); + transformer.applyAdvice( + isMethod() + .and(isPublic()) + .and(named("execute")) + .and(takesArguments(1)) + .and(takesArgument(0, String.class)), + "datadog.trace.instrumentation.postgresql.PgStatementAdvice$ExecuteQueryAdvice"); + transformer.applyAdvice( + isMethod() + .and(isPublic()) + .and(named("addBatch")) + .and(takesArguments(1)) + .and(takesArgument(0, String.class)), + "datadog.trace.instrumentation.postgresql.PgStatementAdvice$AddBatchAdvice"); + transformer.applyAdvice( + isMethod().and(isPublic()).and(named("executeBatch")).and(takesArguments(0)), + "datadog.trace.instrumentation.postgresql.PgStatementAdvice$ExecuteBatchAdvice"); + } +} diff --git a/dd-java-agent/instrumentation/postgresql/postgresql-42.0/src/main/java/datadog/trace/instrumentation/postgresql/PostgreSQLDecorator.java b/dd-java-agent/instrumentation/postgresql/postgresql-42.0/src/main/java/datadog/trace/instrumentation/postgresql/PostgreSQLDecorator.java new file mode 100644 index 00000000000..9d23cbac0c3 --- /dev/null +++ b/dd-java-agent/instrumentation/postgresql/postgresql-42.0/src/main/java/datadog/trace/instrumentation/postgresql/PostgreSQLDecorator.java @@ -0,0 +1,71 @@ +package datadog.trace.instrumentation.postgresql; + +import datadog.trace.api.naming.SpanNaming; +import datadog.trace.bootstrap.instrumentation.api.AgentSpan; +import datadog.trace.bootstrap.instrumentation.api.InternalSpanTypes; +import datadog.trace.bootstrap.instrumentation.api.UTF8BytesString; +import datadog.trace.bootstrap.instrumentation.decorator.DBTypeProcessingDatabaseClientDecorator; +import datadog.trace.bootstrap.instrumentation.jdbc.DBInfo; + +public class PostgreSQLDecorator extends DBTypeProcessingDatabaseClientDecorator { + + public static final PostgreSQLDecorator DECORATE = new PostgreSQLDecorator(); + + public static final CharSequence JAVA_POSTGRESQL = UTF8BytesString.create("java-postgresql"); + public static final CharSequence POSTGRESQL_QUERY = + UTF8BytesString.create( + SpanNaming.instance().namingSchema().database().operation("postgresql")); + private static final String SERVICE_NAME = + SpanNaming.instance().namingSchema().database().service("postgresql"); + + @Override + protected String[] instrumentationNames() { + return new String[] {"postgresql"}; + } + + @Override + protected String service() { + return SERVICE_NAME; + } + + @Override + protected CharSequence component() { + return JAVA_POSTGRESQL; + } + + @Override + protected CharSequence spanType() { + return InternalSpanTypes.SQL; + } + + @Override + protected String dbType() { + return "postgresql"; + } + + @Override + protected String dbUser(final DBInfo info) { + return info.getUser(); + } + + @Override + protected String dbInstance(final DBInfo info) { + if (info.getInstance() != null) { + return info.getInstance(); + } + return info.getDb(); + } + + @Override + protected CharSequence dbHostname(final DBInfo info) { + return info.getHost(); + } + + @Override + protected void postProcessServiceAndOperationName(AgentSpan span, NamingEntry namingEntry) { + if (namingEntry.getService() != null) { + span.setServiceName(namingEntry.getService(), component()); + } + span.setOperationName(namingEntry.getOperation()); + } +} diff --git a/dd-java-agent/instrumentation/postgresql/postgresql-42.0/src/main/java/datadog/trace/instrumentation/postgresql/PostgreSQLModule.java b/dd-java-agent/instrumentation/postgresql/postgresql-42.0/src/main/java/datadog/trace/instrumentation/postgresql/PostgreSQLModule.java new file mode 100644 index 00000000000..4f625028cc6 --- /dev/null +++ b/dd-java-agent/instrumentation/postgresql/postgresql-42.0/src/main/java/datadog/trace/instrumentation/postgresql/PostgreSQLModule.java @@ -0,0 +1,44 @@ +package datadog.trace.instrumentation.postgresql; + +import com.google.auto.service.AutoService; +import datadog.trace.agent.tooling.Instrumenter; +import datadog.trace.agent.tooling.InstrumenterModule; +import java.util.Arrays; +import java.util.Collections; +import java.util.List; +import java.util.Map; + +@AutoService(InstrumenterModule.class) +public final class PostgreSQLModule extends InstrumenterModule.Tracing { + + public PostgreSQLModule() { + super("postgresql"); + } + + @Override + public String[] helperClassNames() { + return new String[] { + packageName + ".PostgreSQLDecorator", + packageName + ".PostgreSQLSQLCommenter", + packageName + ".PgStatementAdvice", + packageName + ".PgStatementAdvice$ExecuteQueryAdvice", + packageName + ".PgStatementAdvice$AddBatchAdvice", + packageName + ".PgStatementAdvice$ExecuteBatchAdvice", + packageName + ".PgPreparedStatementAdvice", + packageName + ".PgPreparedStatementAdvice$ConstructorAdvice", + packageName + ".PgPreparedStatementAdvice$ExecuteAdvice", + }; + } + + @Override + public Map contextStore() { + return Collections.singletonMap( + "java.sql.Statement", "datadog.trace.bootstrap.instrumentation.jdbc.DBInfo"); + } + + @Override + public List typeInstrumentations() { + return Arrays.asList( + new PgStatementInstrumentation(), new PgPreparedStatementInstrumentation()); + } +} diff --git a/dd-java-agent/instrumentation/postgresql/postgresql-42.0/src/main/java/datadog/trace/instrumentation/postgresql/PostgreSQLSQLCommenter.java b/dd-java-agent/instrumentation/postgresql/postgresql-42.0/src/main/java/datadog/trace/instrumentation/postgresql/PostgreSQLSQLCommenter.java new file mode 100644 index 00000000000..a028ebe6df0 --- /dev/null +++ b/dd-java-agent/instrumentation/postgresql/postgresql-42.0/src/main/java/datadog/trace/instrumentation/postgresql/PostgreSQLSQLCommenter.java @@ -0,0 +1,85 @@ +package datadog.trace.instrumentation.postgresql; + +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; +import datadog.trace.bootstrap.instrumentation.jdbc.DBInfo; + +/** + * SQL comment injector for PostgreSQL Database Monitoring (DBM). Prepends a {@code /* ... * /} + * comment to SQL queries containing trace context metadata so the database can correlate queries + * back to traces. + */ +public class PostgreSQLSQLCommenter { + + /** + * Prepends a DBM trace comment to the given SQL string. Returns the original SQL unchanged if DBM + * comment injection is disabled or if the SQL already contains a trace comment. + * + * @param sql the original SQL string + * @param span the active span for trace context + * @param dbInfo the database connection info (may be null) + * @return the SQL string with a prepended trace comment, or the original SQL if injection is not + * applicable + */ + public static String inject(String sql, AgentSpan span, DBInfo dbInfo) { + if (!Config.get().isDbmCommentInjectionEnabled() || sql == null || sql.isEmpty()) { + return sql; + } + + if (span == null) { + return sql; + } + + // Force a sampling decision so traceparent has a valid sampling flag + if (span.forceSamplingDecision() == null) { + return sql; + } + + // Check if the SQL already has a trace comment to avoid duplicate injection + if (hasExistingTraceComment(sql)) { + return sql; + } + + String dbService = span.getServiceName(); + String hostname = dbInfo != null ? dbInfo.getHost() : null; + String dbName = dbInfo != null ? dbInfo.getDb() : null; + + String traceParent = + Config.get().getDbmPropagationMode().equals(DBM_PROPAGATION_MODE_FULL) + ? W3CTraceParent.from(span) + : null; + + String comment = + SharedDBCommenter.buildComment(dbService, "postgresql", hostname, dbName, traceParent); + + if (comment == null || comment.isEmpty()) { + return sql; + } + + // Set the DBM trace injected tag on the span + span.setTag(DBM_TRACE_INJECTED, true); + + // Prepend the comment as a SQL block comment + return "/*" + comment + "*/ " + sql; + } + + /** Checks if the given SQL string already contains a DBM trace comment. */ + private static boolean hasExistingTraceComment(String sql) { + // Look for an opening block comment and check if it contains trace comment markers + int commentStart = sql.indexOf("/*"); + if (commentStart < 0) { + return false; + } + int commentEnd = sql.indexOf("*/", commentStart + 2); + if (commentEnd < 0) { + return false; + } + // Check the comment body for trace comment fields + return SharedDBCommenter.containsTraceComment(sql, commentStart + 2, commentEnd); + } +} diff --git a/dd-java-agent/instrumentation/postgresql/postgresql-42.0/src/test/java/datadog/trace/instrumentation/postgresql/PostgreSQLInstrumentationTest.java b/dd-java-agent/instrumentation/postgresql/postgresql-42.0/src/test/java/datadog/trace/instrumentation/postgresql/PostgreSQLInstrumentationTest.java new file mode 100644 index 00000000000..78b83ff3ee3 --- /dev/null +++ b/dd-java-agent/instrumentation/postgresql/postgresql-42.0/src/test/java/datadog/trace/instrumentation/postgresql/PostgreSQLInstrumentationTest.java @@ -0,0 +1,458 @@ +package datadog.trace.instrumentation.postgresql; + +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.includes; +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_CLIENT_HOST_SPLIT_BY_INSTANCE; +import static datadog.trace.api.config.TraceInstrumentationConfig.DB_DBM_PROPAGATION_MODE_MODE; +import static datadog.trace.bootstrap.instrumentation.api.AgentTracer.activateSpan; +import static datadog.trace.bootstrap.instrumentation.api.AgentTracer.startSpan; +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 datadog.trace.agent.test.AbstractInstrumentationTest; +import datadog.trace.agent.test.assertions.SpanMatcher; +import datadog.trace.agent.test.assertions.TagsMatcher; +import datadog.trace.agent.test.assertions.TraceMatcher; +import datadog.trace.api.Config; +import datadog.trace.api.DDSpanTypes; +import datadog.trace.api.DDTags; +import datadog.trace.api.config.TracerConfig; +import datadog.trace.bootstrap.instrumentation.api.AgentScope; +import datadog.trace.bootstrap.instrumentation.api.AgentSpan; +import datadog.trace.bootstrap.instrumentation.api.InstrumentationTags; +import datadog.trace.bootstrap.instrumentation.api.Tags; +import datadog.trace.test.junit.utils.config.WithConfig; +import datadog.trace.test.junit.utils.config.WithConfigExtension; +import java.sql.Connection; +import java.sql.DriverManager; +import java.sql.PreparedStatement; +import java.sql.ResultSet; +import java.sql.SQLException; +import java.sql.Statement; +import java.time.Duration; +import java.util.ArrayList; +import java.util.List; +import org.junit.jupiter.api.AfterAll; +import org.junit.jupiter.api.BeforeAll; +import org.junit.jupiter.api.Test; +import org.testcontainers.containers.PostgreSQLContainer; + +abstract class PostgreSQLInstrumentationTest extends AbstractInstrumentationTest { + + static PostgreSQLContainer container; + static int port; + static Connection connection; + + abstract String service(); + + abstract String operation(); + + @BeforeAll + static void setupAll() throws Exception { + container = + new PostgreSQLContainer<>("postgres:13-alpine") + .withDatabaseName("testdb") + .withUsername("testuser") + .withPassword("testpass") + .withStartupTimeout(Duration.ofMinutes(3)); + container.start(); + port = container.getMappedPort(PostgreSQLContainer.POSTGRESQL_PORT); + + // Connect with retries to handle port forwarding delays (e.g. Colima) + int maxRetries = 10; + for (int i = 0; i < maxRetries; i++) { + try { + connection = + DriverManager.getConnection( + container.getJdbcUrl(), container.getUsername(), container.getPassword()); + break; + } catch (SQLException e) { + if (i == maxRetries - 1) { + throw e; + } + Thread.sleep(2000); + } + } + + // Set up the test schema + Statement stmt = connection.createStatement(); + stmt.execute( + "CREATE TABLE IF NOT EXISTS test_table (id SERIAL PRIMARY KEY, name VARCHAR(100))"); + stmt.execute("INSERT INTO test_table (name) VALUES ('alice')"); + stmt.execute("INSERT INTO test_table (name) VALUES ('bob')"); + stmt.close(); + + // Allow time for any setup traces to be written, then clear for actual tests + Thread.sleep(1000); + writer.start(); + } + + @AfterAll + static void cleanupAll() throws Exception { + if (connection != null) { + connection.close(); + } + if (container != null) { + container.stop(); + } + } + + @Test + void executeQueryCreatesSpan() throws Exception { + Statement stmt = connection.createStatement(); + try { + ResultSet rs = stmt.executeQuery("SELECT * FROM test_table"); + rs.next(); + rs.close(); + } finally { + stmt.close(); + } + + assertTraces(trace(postgresSpan("SELECT * FROM test_table", "SELECT", false, false, false))); + } + + @Test + void executeUpdateCreatesSpan() throws Exception { + Statement stmt = connection.createStatement(); + try { + int count = stmt.executeUpdate("INSERT INTO test_table (name) VALUES ('charlie')"); + assert count == 1; + } finally { + stmt.close(); + } + + assertTraces( + trace( + postgresSpan( + "INSERT INTO test_table (name) VALUES (?)", "INSERT", false, false, false))); + } + + @Test + void executeCreatesSpan() throws Exception { + Statement stmt = connection.createStatement(); + try { + stmt.execute("SELECT 1"); + } finally { + stmt.close(); + } + + assertTraces(trace(postgresSpan("SELECT ?", "SELECT", false, false, false))); + } + + @Test + void executeBatchCreatesSpan() throws Exception { + Statement stmt = connection.createStatement(); + try { + stmt.addBatch("INSERT INTO test_table (name) VALUES ('batch1')"); + stmt.addBatch("INSERT INTO test_table (name) VALUES ('batch2')"); + int[] results = stmt.executeBatch(); + assert results.length == 2; + } finally { + stmt.close(); + } + + assertTraces( + trace( + postgresSpan( + "INSERT INTO test_table (name) VALUES (?)", "INSERT", false, false, false))); + } + + @Test + void preparedStatementExecuteQueryCreatesSpan() throws Exception { + PreparedStatement pstmt = + connection.prepareStatement("SELECT * FROM test_table WHERE name = ?"); + try { + pstmt.setString(1, "alice"); + ResultSet rs = pstmt.executeQuery(); + rs.next(); + assert "alice".equals(rs.getString("name")); + rs.close(); + } finally { + pstmt.close(); + } + + assertTraces( + trace( + postgresSpan( + "SELECT * FROM test_table WHERE name = ?", "SELECT", false, false, false))); + } + + @Test + void preparedStatementExecuteUpdateCreatesSpan() throws Exception { + PreparedStatement pstmt = + connection.prepareStatement("INSERT INTO test_table (name) VALUES (?)"); + try { + pstmt.setString(1, "prepared_insert"); + int count = pstmt.executeUpdate(); + assert count == 1; + } finally { + pstmt.close(); + } + + assertTraces( + trace( + postgresSpan( + "INSERT INTO test_table (name) VALUES (?)", "INSERT", false, false, false))); + } + + @Test + void preparedStatementExecuteCreatesSpan() throws Exception { + PreparedStatement pstmt = connection.prepareStatement("SELECT * FROM test_table WHERE id = ?"); + try { + pstmt.setInt(1, 1); + pstmt.execute(); + } finally { + pstmt.close(); + } + + assertTraces( + trace( + postgresSpan("SELECT * FROM test_table WHERE id = ?", "SELECT", false, false, false))); + } + + @Test + void preparedStatementExecuteBatchCreatesSpan() throws Exception { + PreparedStatement pstmt = + connection.prepareStatement("INSERT INTO test_table (name) VALUES (?)"); + try { + pstmt.setString(1, "batch_prepared1"); + pstmt.addBatch(); + pstmt.setString(1, "batch_prepared2"); + pstmt.addBatch(); + int[] results = pstmt.executeBatch(); + assert results.length == 2; + } finally { + pstmt.close(); + } + + assertTraces( + trace( + postgresSpan( + "INSERT INTO test_table (name) VALUES (?)", "INSERT", false, false, false))); + } + + @Test + void queryUnderParentSpan() throws Exception { + AgentSpan parentSpan = startSpan("test", "parent"); + AgentScope scope = activateSpan(parentSpan); + try { + Statement stmt = connection.createStatement(); + stmt.executeQuery("SELECT * FROM test_table"); + stmt.close(); + } finally { + scope.close(); + parentSpan.finish(); + } + + assertTraces( + trace( + TraceMatcher.SORT_BY_START_TIME, + SpanMatcher.span() + .serviceNameDefined() + .operationName("parent") + .resourceName("parent") + .root(), + postgresSpan("SELECT * FROM test_table", "SELECT", false, false, false) + .childOfPrevious())); + } + + @Test + void errorDuringQuery() throws Exception { + Statement stmt = connection.createStatement(); + try { + stmt.executeQuery("SELECT * FROM non_existent_table"); + } catch (SQLException expected) { + // expected + } finally { + stmt.close(); + } + + assertTraces( + trace(postgresSpan("SELECT * FROM non_existent_table", "SELECT", true, false, false))); + } + + @Test + void errorDuringExecuteUpdate() throws Exception { + Statement stmt = connection.createStatement(); + try { + stmt.executeUpdate("INSERT INTO non_existent_table (name) VALUES ('fail')"); + } catch (SQLException expected) { + // expected + } finally { + stmt.close(); + } + + assertTraces( + trace( + postgresSpan( + "INSERT INTO non_existent_table (name) VALUES (?)", "INSERT", true, false, false))); + } + + @Test + void splitByInstance() throws Exception { + WithConfigExtension.injectSysConfig(DB_CLIENT_HOST_SPLIT_BY_INSTANCE, "true"); + Statement stmt = connection.createStatement(); + try { + stmt.executeQuery("SELECT * FROM test_table"); + } finally { + stmt.close(); + } + + assertTraces(trace(postgresSpan("SELECT * FROM test_table", "SELECT", false, true, false))); + } + + @Test + void dbmServiceModeInjectsCommentAndSetsTag() throws Exception { + WithConfigExtension.injectSysConfig(DB_DBM_PROPAGATION_MODE_MODE, "service"); + Statement stmt = connection.createStatement(); + try { + ResultSet rs = stmt.executeQuery("SELECT * FROM test_table"); + rs.next(); + rs.close(); + } finally { + stmt.close(); + } + + assertTraces(trace(postgresSpan("SELECT * FROM test_table", "SELECT", false, false, true))); + } + + @Test + void dbmFullModeInjectsCommentWithTraceparent() throws Exception { + WithConfigExtension.injectSysConfig(DB_DBM_PROPAGATION_MODE_MODE, "full"); + Statement stmt = connection.createStatement(); + try { + ResultSet rs = stmt.executeQuery("SELECT * FROM test_table"); + rs.next(); + rs.close(); + } finally { + stmt.close(); + } + + assertTraces(trace(postgresSpan("SELECT * FROM test_table", "SELECT", false, false, true))); + } + + @Test + void dbmServiceModeInjectsCommentForPreparedStatement() throws Exception { + WithConfigExtension.injectSysConfig(DB_DBM_PROPAGATION_MODE_MODE, "service"); + PreparedStatement pstmt = + connection.prepareStatement("SELECT * FROM test_table WHERE name = ?"); + try { + pstmt.setString(1, "alice"); + ResultSet rs = pstmt.executeQuery(); + rs.next(); + assert "alice".equals(rs.getString("name")); + rs.close(); + } finally { + pstmt.close(); + } + + assertTraces( + trace( + postgresSpan("SELECT * FROM test_table WHERE name = ?", "SELECT", false, false, true))); + } + + @Test + void dbmFullModeInjectsCommentForPreparedStatement() throws Exception { + WithConfigExtension.injectSysConfig(DB_DBM_PROPAGATION_MODE_MODE, "full"); + PreparedStatement pstmt = + connection.prepareStatement("SELECT * FROM test_table WHERE name = ?"); + try { + pstmt.setString(1, "alice"); + ResultSet rs = pstmt.executeQuery(); + rs.next(); + assert "alice".equals(rs.getString("name")); + rs.close(); + } finally { + pstmt.close(); + } + + assertTraces( + trace( + postgresSpan("SELECT * FROM test_table WHERE name = ?", "SELECT", false, false, true))); + } + + @Test + void dbmDisabledDoesNotInjectCommentOrSetTag() throws Exception { + WithConfigExtension.injectSysConfig(DB_DBM_PROPAGATION_MODE_MODE, "disabled"); + Statement stmt = connection.createStatement(); + try { + ResultSet rs = stmt.executeQuery("SELECT * FROM test_table"); + rs.next(); + rs.close(); + } finally { + stmt.close(); + } + + assertTraces(trace(postgresSpan("SELECT * FROM test_table", "SELECT", false, false, false))); + } + + SpanMatcher postgresSpan( + String resource, + String dbOperation, + boolean hasError, + boolean renameService, + boolean addDbmTag) { + List tagMatchers = new ArrayList<>(); + tagMatchers.add(defaultTags()); + tagMatchers.add(tag(Tags.COMPONENT, is("java-postgresql"))); + tagMatchers.add(tag(Tags.SPAN_KIND, is(Tags.SPAN_KIND_CLIENT))); + tagMatchers.add(tag(Tags.DB_TYPE, is("postgresql"))); + tagMatchers.add(tag(Tags.DB_INSTANCE, is("testdb"))); + tagMatchers.add(tag(Tags.DB_USER, is("testuser"))); + tagMatchers.add(tag(Tags.DB_OPERATION, is(dbOperation))); + tagMatchers.add(tag(Tags.PEER_HOSTNAME, isNonNull())); + tagMatchers.add(tag(Tags.PEER_PORT, isNonNull())); + if (hasError) { + tagMatchers.add(error(SQLException.class)); + tagMatchers.add(tag(DDTags.ERROR_MSG, isNonNull())); + } + if (addDbmTag) { + tagMatchers.add(tag(InstrumentationTags.DBM_TRACE_INJECTED, is(true))); + } + // Peer service tags + tagMatchers.add(includes(DDTags.PEER_SERVICE_SOURCE)); + tagMatchers.add(tag("peer.service", any())); + tagMatchers.add(tag(DDTags.DD_SVC_SRC, any())); + + return SpanMatcher.span() + .serviceName(renameService ? "testdb" : service()) + .operationName(operation()) + .resourceName(resource) + .type(DDSpanTypes.SQL) + .error(hasError) + .root() + .tags(tagMatchers.toArray(new TagsMatcher[0])); + } +} + +@WithConfig(key = TracerConfig.TRACE_SPAN_ATTRIBUTE_SCHEMA, value = "v0") +class PostgreSQLInstrumentationV0Test extends PostgreSQLInstrumentationTest { + + @Override + String service() { + return "postgresql"; + } + + @Override + String operation() { + return "postgresql.query"; + } +} + +@WithConfig(key = TracerConfig.TRACE_SPAN_ATTRIBUTE_SCHEMA, value = "v1") +class PostgreSQLInstrumentationV1ForkedTest extends PostgreSQLInstrumentationTest { + + @Override + String service() { + return Config.get().getServiceName(); + } + + @Override + String operation() { + return "postgresql.query"; + } +} diff --git a/dd-java-agent/instrumentation/postgresql/postgresql-42.0/src/test/resources/pg_hba.conf b/dd-java-agent/instrumentation/postgresql/postgresql-42.0/src/test/resources/pg_hba.conf new file mode 100644 index 00000000000..6b4940dcf8d --- /dev/null +++ b/dd-java-agent/instrumentation/postgresql/postgresql-42.0/src/test/resources/pg_hba.conf @@ -0,0 +1,4 @@ +# TYPE DATABASE USER ADDRESS METHOD +local all all md5 +host all all 0.0.0.0/0 md5 +host all all ::/0 md5 diff --git a/metadata/agent-jar-checks.properties b/metadata/agent-jar-checks.properties index fd04d637e35..bfe0b4ae98d 100644 --- a/metadata/agent-jar-checks.properties +++ b/metadata/agent-jar-checks.properties @@ -142,6 +142,7 @@ expected.integrations = IastInstrumentation,\ pekko_actor_send,\ play,\ play-ws,\ + postgresql,\ protobuf,\ quartz,\ ratpack,\ diff --git a/metadata/supported-configurations.json b/metadata/supported-configurations.json index 6241dc3e6af..d677e45cfe8 100644 --- a/metadata/supported-configurations.json +++ b/metadata/supported-configurations.json @@ -8961,6 +8961,30 @@ "aliases": ["DD_TRACE_INTEGRATION_PLAY_WS_ENABLED", "DD_INTEGRATION_PLAY_WS_ENABLED"] } ], + "DD_TRACE_POSTGRESQL_ANALYTICS_ENABLED": [ + { + "version": "A", + "type": "boolean", + "default": "false", + "aliases": ["DD_POSTGRESQL_ANALYTICS_ENABLED"] + } + ], + "DD_TRACE_POSTGRESQL_ANALYTICS_SAMPLE_RATE": [ + { + "version": "A", + "type": "decimal", + "default": "1.0", + "aliases": ["DD_POSTGRESQL_ANALYTICS_SAMPLE_RATE"] + } + ], + "DD_TRACE_POSTGRESQL_ENABLED": [ + { + "version": "A", + "type": "boolean", + "default": "true", + "aliases": ["DD_TRACE_INTEGRATION_POSTGRESQL_ENABLED", "DD_INTEGRATION_POSTGRESQL_ENABLED"] + } + ], "DD_TRACE_POST_PROCESSING_TIMEOUT": [ { "version": "A", diff --git a/settings.gradle.kts b/settings.gradle.kts index 86aa23cc9d4..2c172cf1d4e 100644 --- a/settings.gradle.kts +++ b/settings.gradle.kts @@ -417,7 +417,7 @@ include( ":dd-java-agent:instrumentation:javax-xml-1.4", ":dd-java-agent:instrumentation:jboss:jboss-logmanager-1.1", ":dd-java-agent:instrumentation:jboss:jboss-modules-1.3", - ":dd-java-agent:instrumentation:jdbc:scalikejdbc-3.5", + // ":dd-java-agent:instrumentation:jdbc:scalikejdbc-3.5", // Commented out - directory doesn't exist ":dd-java-agent:instrumentation:jdbc", ":dd-java-agent:instrumentation:jedis:jedis-1.4", ":dd-java-agent:instrumentation:jedis:jedis-3.0", @@ -535,6 +535,7 @@ include( ":dd-java-agent:instrumentation:play:play-appsec-2.6", ":dd-java-agent:instrumentation:play:play-appsec-2.7", ":dd-java-agent:instrumentation:play:play-appsec-common", + ":dd-java-agent:instrumentation:postgresql:postgresql-42.0", ":dd-java-agent:instrumentation:protobuf-3.0", ":dd-java-agent:instrumentation:quartz-2.0", ":dd-java-agent:instrumentation:rabbitmq-amqp-2.7", diff --git a/utils/test-junit-utils/src/main/java/datadog/trace/test/junit/utils/assertions/Is.java b/utils/test-junit-utils/src/main/java/datadog/trace/test/junit/utils/assertions/Is.java index f7a3346e64e..f05cffadbdd 100644 --- a/utils/test-junit-utils/src/main/java/datadog/trace/test/junit/utils/assertions/Is.java +++ b/utils/test-junit-utils/src/main/java/datadog/trace/test/junit/utils/assertions/Is.java @@ -27,6 +27,13 @@ public String failureReason() { @Override public boolean test(T t) { - return this.expected.equals(t); + if (this.expected.equals(t)) { + return true; + } + // Handle CharSequence comparison (e.g. String vs UTF8BytesString) + if (this.expected instanceof CharSequence && t instanceof CharSequence) { + return this.expected.toString().equals(t.toString()); + } + return false; } }