diff --git a/.agents/skills/apm-integrations/references/instrumenter-module.md b/.agents/skills/apm-integrations/references/instrumenter-module.md index 48059ed8f1d..d138970f626 100644 --- a/.agents/skills/apm-integrations/references/instrumenter-module.md +++ b/.agents/skills/apm-integrations/references/instrumenter-module.md @@ -41,6 +41,78 @@ ✅ `public String[] triggerClasses() { return new String[]{"com.example.Foo"}; }` +### Before hand-writing method advice, check for an official interception SPI + +Before instrumenting a client library by matching and advising its methods directly, check whether the library — or an official companion library in its ecosystem — exposes a first-class **interception / listener SPI**: a listener interface, an interceptor-registration hook, or a proxy/wrapper factory built for observability. Many client libraries provide one, and it is almost always the better hook point than hand-written method advice. + +Examples of such SPIs: + +- **R2DBC** → `r2dbc-proxy` (an official R2DBC project) exposes `ProxyExecutionListener` / `ProxyMethodExecutionListener`, invoked around every query via `beforeQuery(QueryExecutionInfo)` / `afterQuery(QueryExecutionInfo)`. You install it by wrapping the `ConnectionFactory`. +- **gRPC** → `ClientInterceptor` / `ServerInterceptor`. +- **Kafka** → `ProducerInterceptor` / `ConsumerInterceptor` (`interceptor.classes`). +- **JAX-RS / JDBC / many others** → filter, interceptor, or proxy-driver registration hooks. + +**Why prefer the SPI over hand-written advice, especially for reactive/async clients:** a listener SPI is designed around the library's own execution lifecycle and delivers a single well-defined start/end (and error/**cancel**) callback per logical operation, with the operation's metadata (query text, connection info, status, timing) already assembled for you. Hand-written advice on a reactive method must instead re-implement that lifecycle by wrapping the returned `Publisher`/`Mono`/`Flux` and tracking subscribe/complete/error/**cancel** by hand — this is easy to get subtly wrong. A wrapper that finishes the span only on `onComplete`/`onError` will **leak every span that gets cancelled** (reactive pipelines routinely cancel upstreams — e.g. `take(1)`, timeouts, `DiscardOnCancel` operators), so the span is created but never finished and never exported. It is also easy to wrap at the wrong granularity and emit many spans per logical query. The library's listener already handles all of this. + +**Do this without forcing a new runtime dependency on the user.** Using a library's listener SPI does NOT require adding it as a normal (`implementation`/`api`) dependency: + +- Declare the interception library **`compileOnly`** so it is not put on the application's runtime classpath, and inject your listener implementation and any glue via the module's `helperClassNames()` (the same mechanism used for decorators and other injected helpers). The listener classes travel inside the agent, not the user's app. +- If helper injection is impractical for a given SPI, the alternative is to **shade/bundle** the interception library's classes into the instrumentation rather than depend on it at runtime. + +Reach for hand-written method advice when no such SPI exists, or when the SPI cannot express what you need to capture. When one does exist and fits, prefer it — and note the choice (and the `compileOnly`/injection approach) in the PR so a reviewer sees the dependency was considered. + +**Worked example — hooking a registration/factory point instead of the query methods (R2DBC + `r2dbc-proxy`):** + +1. **Advice target is the factory, not the query.** Hook `io.r2dbc.spi.ConnectionFactories.find(ConnectionFactoryOptions)` — NOT `Statement.execute()` / `Batch.execute()`. Replace the returned `ConnectionFactory` using `@Advice.AssignReturned.ToReturned` on method exit: + + ```java + @Advice.OnMethodExit(suppress = Throwable.class, inline = false) + @Advice.AssignReturned.ToReturned + public static ConnectionFactory wrap( + @Advice.Return ConnectionFactory factory, + @Advice.Argument(0) ConnectionFactoryOptions options) { + return R2dbcTracingSupport.wrapConnectionFactory(factory, options); + } + ``` + + Do NOT use a writable `@Advice.Return(readOnly = false) Publisher<...> ...` here — binding a reactive-streams interface to a writable return against a concrete `Mono`/`Flux` result fails Byte Buddy transformation. `@Advice.AssignReturned.ToReturned` sidesteps this because it substitutes the value rather than mutating a typed slot. + +2. **Wrap helper installs the listener** (injected via `helperClassNames()`, `r2dbc-proxy` declared `compileOnly`): + + ```java + public static ConnectionFactory wrapConnectionFactory( + ConnectionFactory factory, ConnectionFactoryOptions options) { + ProxyConfig cfg = new ProxyConfig(); + cfg.addListener(new TraceProxyListener(options)); + return ProxyConnectionFactory.builder(factory, cfg).build(); + } + ``` + +3. **Listener drives the span lifecycle** — implements `ProxyMethodExecutionListener`; stash the span on the query's own value store (the interception library already carries one per call — do not add a separate `ContextStore` keyed on the driver object): + + ```java + public class TraceProxyListener implements ProxyMethodExecutionListener { + @Override + public void beforeQuery(QueryExecutionInfo qei) { + AgentSpan span = startSpan(...); + DECORATE.afterStart(span); + DECORATE.onStatement(span, qei); + qei.getValueStore().put(SPAN_KEY, span); + } + @Override + public void afterQuery(QueryExecutionInfo qei) { + AgentSpan span = qei.getValueStore().get(SPAN_KEY, AgentSpan.class); + if (qei.getThrowable() != null) DECORATE.onError(span, qei.getThrowable()); + DECORATE.beforeFinish(span); + span.finish(); + } + } + ``` + + `beforeQuery`/`afterQuery` fire once per logical operation and `r2dbc-proxy` itself owns completion/error/**cancel** — you do not hand-roll a `Publisher` wrapper to catch those. + +Net result: 1 advice (factory hook) + 1 wrap helper + 1 listener class — smaller than the method-advice alternative (which needs per-method advice on both `Statement.execute()` and `Batch.execute()`, plus a hand-rolled cancel-safe `Publisher` wrapper) and correct by construction on cancellation. The same shape applies to any other SPI in the examples list above: hook the registration/interceptor-installation point, not the client's own operational methods. + ### Before writing a new module, scan for an existing one Before creating `dd-java-agent/instrumentation/$framework/$framework-$version/`, check whether `dd-java-agent/instrumentation/$framework/` already exists and what's in it. @@ -123,7 +195,7 @@ A helper class is appropriate when multiple instrumentation classes share the sa For database-client integrations (`DatabaseClientDecorator` / `DBTypeProcessingDatabaseClientDecorator`), capture connection metadata (host, port, db name, user) at **connection-establishment** time and cache it in a `ContextStore` keyed on the connection object — not lazily on the first query. The canonical pattern is a dedicated instrumentation on the connect/factory method: - **JDBC** — `dd-java-agent/instrumentation/jdbc/DriverInstrumentation.java` hooks `Driver.connect(url, props)` and populates `InstrumentationContext.get(Connection.class, DBInfo.class)` at open time. Statement advice then reads the already-cached `DBInfo`. -- **Reactive drivers with an async connect** — the equivalent connect point is the connection FACTORY, not the connection object. For R2DBC, `io.r2dbc.spi.ConnectionFactoryOptions` is the only place host/port/database/user are exposed as structured data; `io.r2dbc.spi.ConnectionMetadata` (on the live `Connection`) exposes ONLY product name/version. **But `ConnectionFactory.create()` is a zero-argument SPI method returning a `Publisher` — the options are NOT available at `create()`.** Capture them earlier, where the factory is built: hook `ConnectionFactories.get(ConnectionFactoryOptions)` (or the provider-construction path) and store the options in a `ContextStore`; then, in advice on `create()`, read the stored options for that factory and thread them onto the asynchronously-emitted `Connection` (a second context store keyed on the returned `Connection`). Hooking only `Connection.createStatement()` + `ConnectionMetadata` CANNOT populate `db.name`/`peer.hostname`/`db.user`/port. (OpenTelemetry's R2DBC instrumentation does exactly this options→factory→connection threading; it is a good reference.) +- **Reactive drivers with an async connect** — the equivalent connect point is the connection FACTORY, not the connection object. For R2DBC, `io.r2dbc.spi.ConnectionFactoryOptions` is the only place host/port/database/user are exposed as structured data; `io.r2dbc.spi.ConnectionMetadata` (on the live `Connection`) exposes ONLY product name/version. **But `ConnectionFactory.create()` is a zero-argument SPI method returning a `Publisher` — the options are NOT available at `create()`.** Capture them earlier, where the factory is built: hook `ConnectionFactories.get(ConnectionFactoryOptions)` (or the provider-construction path) and store the options in a `ContextStore`; then, in advice on `create()`, read the stored options for that factory and thread them onto the asynchronously-emitted `Connection` (a second context store keyed on the returned `Connection`). Hooking only `Connection.createStatement()` + `ConnectionMetadata` CANNOT populate `db.name`/`peer.hostname`/`db.user`/port. (OpenTelemetry's R2DBC instrumentation does exactly this options→factory→connection threading; it is a good reference.) Note also — per "Before hand-writing method advice, check for an official interception SPI" above — that R2DBC has `r2dbc-proxy`, whose `ProxyMethodExecutionListener` surfaces this connection/query metadata through a listener and handles the reactive lifecycle for you; prefer it over hand-wrapping the `create()`/`execute()` publishers where it fits. Why eager-at-connect beats lazy-per-query: lazy extraction (e.g. `statement.getConnection().getMetaData().getURL()` on first execute) works for plain JDBC but (a) pays the extraction cost on every connection's first query instead of amortizing at pool-open, and (b) silently yields nothing when the metadata is not reachable from the object the query advice happens to hold — which is exactly what happens for reactive drivers whose statement/connection objects don't carry the factory options. diff --git a/dd-java-agent/instrumentation/build.gradle b/dd-java-agent/instrumentation/build.gradle index 924b627e8b6..6d318f7a5fe 100644 --- a/dd-java-agent/instrumentation/build.gradle +++ b/dd-java-agent/instrumentation/build.gradle @@ -135,6 +135,10 @@ tasks.named('shadowJar', ShadowJar) { exclude(dependency('com.google.re2j:re2j')) deps.excludeShared.execute(it) } + // Redundant metadata file duplicated across sibling io.r2dbc:* artifacts + // (r2dbc-spi, r2dbc-proxy) once a module bundles more than one of them. + // Not required content — drop it instead of failing the aggregate jar. + exclude 'META-INF/CHANGELOG' } // temporary config to add slf4j-simple so we get logging from instrumenters while indexing diff --git a/dd-java-agent/instrumentation/r2dbc/r2dbc-1.0/build.gradle b/dd-java-agent/instrumentation/r2dbc/r2dbc-1.0/build.gradle new file mode 100644 index 00000000000..a2950b3f3d7 --- /dev/null +++ b/dd-java-agent/instrumentation/r2dbc/r2dbc-1.0/build.gradle @@ -0,0 +1,42 @@ +plugins { + id 'dd-trace-java.module.instrumentation' +} + +muzzle { + pass { + group = "io.r2dbc" + module = "r2dbc-spi" + versions = "[1.0.0.RELEASE,)" + // r2dbc-proxy is compileOnly and referenced directly by the injected + // listener/wrap-helper classes (R2dbcTracingSupport, TraceProxyExecutionListener). + // Muzzle only puts the pinned primary dependency (r2dbc-spi) on the test + // classpath by default, so the compileOnly interception library must be + // added explicitly or muzzle reports its classes as "missing". + extraDependency 'io.r2dbc:r2dbc-proxy:1.1.0.RELEASE' + } +} + +addTestSuiteForDir('latestDepTest', 'test') + +dependencies { + compileOnly group: 'io.r2dbc', name: 'r2dbc-spi', version: '1.0.0.RELEASE' + // r2dbc-proxy must be bundled into the agent (implementation, NOT compileOnly): + // R2dbcTracingSupport/TraceProxyExecutionListener reference its types directly, and + // ordinary R2DBC apps do not depend on r2dbc-proxy themselves. compileOnly would leave + // those classes absent from the shaded agent jar, and the runtime muzzle safety check + // then blocks the whole module to avoid a NoClassDefFoundError once the target app's + // classloader is checked. r2dbc-spi stays compileOnly — target apps DO provide that one. + implementation group: 'io.r2dbc', name: 'r2dbc-proxy', version: '1.1.0.RELEASE' + + testImplementation group: 'io.r2dbc', name: 'r2dbc-spi', version: '1.0.0.RELEASE' + testImplementation group: 'io.r2dbc', name: 'r2dbc-proxy', version: '1.1.0.RELEASE' + // H2 R2DBC driver for in-memory database testing + testImplementation group: 'io.r2dbc', name: 'r2dbc-h2', version: '1.0.0.RELEASE' + // Reactor for blocking on reactive types in tests + testImplementation group: 'io.projectreactor', name: 'reactor-core', version: '3.5.0' + + latestDepTestImplementation group: 'io.r2dbc', name: 'r2dbc-spi', version: '1.+' + latestDepTestImplementation group: 'io.r2dbc', name: 'r2dbc-proxy', version: '1.+' + latestDepTestImplementation group: 'io.r2dbc', name: 'r2dbc-h2', version: '1.+' + latestDepTestImplementation group: 'io.projectreactor', name: 'reactor-core', version: '3.+' +} diff --git a/dd-java-agent/instrumentation/r2dbc/r2dbc-1.0/src/main/java/datadog/trace/instrumentation/r2dbc/R2dbcConnectionCallbackInstrumentation.java b/dd-java-agent/instrumentation/r2dbc/r2dbc-1.0/src/main/java/datadog/trace/instrumentation/r2dbc/R2dbcConnectionCallbackInstrumentation.java new file mode 100644 index 00000000000..61104e07477 --- /dev/null +++ b/dd-java-agent/instrumentation/r2dbc/r2dbc-1.0/src/main/java/datadog/trace/instrumentation/r2dbc/R2dbcConnectionCallbackInstrumentation.java @@ -0,0 +1,207 @@ +package datadog.trace.instrumentation.r2dbc; + +import static datadog.trace.agent.tooling.bytebuddy.matcher.NameMatchers.named; +import static datadog.trace.instrumentation.r2dbc.R2dbcDecorator.DECORATE; +import static net.bytebuddy.matcher.ElementMatchers.isMethod; +import static net.bytebuddy.matcher.ElementMatchers.takesArgument; +import static net.bytebuddy.matcher.ElementMatchers.takesArguments; + +import com.google.auto.service.AutoService; +import datadog.trace.agent.tooling.Instrumenter; +import datadog.trace.agent.tooling.InstrumenterModule; +import datadog.trace.api.Config; +import io.r2dbc.proxy.core.ConnectionInfo; +import io.r2dbc.spi.ConnectionFactoryOptions; +import java.lang.reflect.Method; +import net.bytebuddy.asm.Advice; + +/** + * Instruments {@code io.r2dbc.proxy.callback.ConnectionCallbackHandler} to inject DBM SQL comments + * into queries before they reach the database driver. This is the R2DBC equivalent of JDBC's {@code + * DBMCompatibleConnectionInstrumentation}. + * + *

The r2dbc-proxy library uses JDK dynamic proxies for Connection objects, so we cannot + * instrument them with ByteBuddy directly. Instead, we intercept the callback handler's {@code + * invoke} method which is called for every method on the proxied Connection. When {@code + * createStatement(String)} is invoked, we inject the SQL comment into the first argument. + */ +@AutoService(InstrumenterModule.class) +public class R2dbcConnectionCallbackInstrumentation extends InstrumenterModule.Tracing + implements Instrumenter.ForSingleType, Instrumenter.HasMethodAdvice { + + public R2dbcConnectionCallbackInstrumentation() { + super("r2dbc"); + } + + @Override + public String instrumentedType() { + return "io.r2dbc.proxy.callback.ConnectionCallbackHandler"; + } + + @Override + public String[] helperClassNames() { + return new String[] { + // See R2dbcInstrumentation#helperClassNames for why the full r2dbc-proxy class set + // (rather than a hand-picked subset) is required, and for the topological ordering + // rationale (supertypes must be injected before their implementing classes). + "io.r2dbc.proxy.callback.AfterQueryCallbackInvoker", + "io.r2dbc.proxy.callback.CallbackHandler", + "io.r2dbc.proxy.callback.CallbackHandlerSupport", + "io.r2dbc.proxy.callback.BatchCallbackHandler", + "io.r2dbc.proxy.callback.CallbackHandlerSupport$MethodInvocationStrategy", + "io.r2dbc.proxy.callback.ConnectionCallbackHandler", + "io.r2dbc.proxy.callback.ConnectionFactoryCallbackHandler", + "io.r2dbc.proxy.callback.MethodInvocationSubscriber", + "io.r2dbc.proxy.callback.ConnectionFactoryCreateMethodInvocationSubscriber", + "io.r2dbc.proxy.callback.ConnectionHolder", + "io.r2dbc.proxy.callback.ConnectionIdManager", + "io.r2dbc.proxy.callback.DefaultConnectionIdManager", + "io.r2dbc.proxy.core.ConnectionInfo", + "io.r2dbc.proxy.callback.DefaultConnectionInfo", + "io.r2dbc.proxy.callback.DelegatingContextView", + "io.r2dbc.proxy.callback.ProxyFactory", + "io.r2dbc.proxy.callback.JdkProxyFactory", + "io.r2dbc.proxy.callback.JdkProxyFactory$CallbackInvocationHandler", + "io.r2dbc.proxy.callback.ProxyFactoryFactory", + "io.r2dbc.proxy.callback.JdkProxyFactoryFactory", + "io.r2dbc.proxy.core.BindInfo", + "io.r2dbc.proxy.callback.MutableBindInfo", + "io.r2dbc.proxy.core.MethodExecutionInfo", + "io.r2dbc.proxy.callback.MutableMethodExecutionInfo", + "io.r2dbc.proxy.core.QueryExecutionInfo", + "io.r2dbc.proxy.callback.MutableQueryExecutionInfo", + "io.r2dbc.proxy.core.StatementInfo", + "io.r2dbc.proxy.callback.MutableStatementInfo", + "io.r2dbc.proxy.callback.ProxyConfig", + "io.r2dbc.proxy.callback.ProxyConfig$1", + "io.r2dbc.proxy.callback.ProxyConfig$Builder", + "io.r2dbc.proxy.callback.ProxyConfigHolder", + "io.r2dbc.proxy.callback.ProxyUtils", + "io.r2dbc.proxy.callback.QueriesExecutionContext", + "io.r2dbc.proxy.callback.QueryInvocationSubscriber", + "io.r2dbc.proxy.callback.ResultCallbackHandler", + "io.r2dbc.proxy.callback.ResultInvocationSubscriber", + "io.r2dbc.proxy.callback.RowCallbackHandler", + "io.r2dbc.proxy.callback.StatementCallbackHandler", + "io.r2dbc.proxy.callback.StopWatch", + "io.r2dbc.proxy.core.Binding", + "io.r2dbc.proxy.core.Bindings", + "io.r2dbc.proxy.core.Bindings$1", + "io.r2dbc.proxy.core.Bindings$IndexBinding", + "io.r2dbc.proxy.core.Bindings$NamedBinding", + "io.r2dbc.proxy.core.BoundValue", + "io.r2dbc.proxy.core.BoundValue$DefaultBoundValue", + "io.r2dbc.proxy.core.ValueStore", + "io.r2dbc.proxy.core.DefaultValueStore", + "io.r2dbc.proxy.core.ExecutionType", + "io.r2dbc.proxy.core.ProxyEventType", + "io.r2dbc.proxy.core.QueryInfo", + "io.r2dbc.proxy.core.R2dbcProxyException", + "io.r2dbc.proxy.listener.BindParameterConverter", + "io.r2dbc.proxy.listener.BindParameterConverter$1", + "io.r2dbc.proxy.listener.BindParameterConverter$BindOperation", + "io.r2dbc.proxy.listener.ProxyExecutionListener", + "io.r2dbc.proxy.listener.CompositeProxyExecutionListener", + "io.r2dbc.proxy.listener.LastExecutionAwareListener", + "io.r2dbc.proxy.listener.ProxyMethodExecutionListener", + "io.r2dbc.proxy.listener.ProxyMethodExecutionListenerAdapter", + "io.r2dbc.proxy.listener.ResultRowConverter", + "io.r2dbc.proxy.listener.ResultRowConverter$GetOperation", + "io.r2dbc.proxy.ProxyConnectionFactory", + "io.r2dbc.proxy.ProxyConnectionFactory$1", + "io.r2dbc.proxy.ProxyConnectionFactory$Builder", + "io.r2dbc.proxy.ProxyConnectionFactory$Builder$1", + "io.r2dbc.proxy.ProxyConnectionFactory$Builder$2", + "io.r2dbc.proxy.ProxyConnectionFactory$Builder$3", + "io.r2dbc.proxy.ProxyConnectionFactory$Builder$4", + "io.r2dbc.proxy.ProxyConnectionFactory$Builder$5", + "io.r2dbc.proxy.ProxyConnectionFactoryProvider", + "io.r2dbc.proxy.support.FormatterUtils", + "io.r2dbc.proxy.support.MethodExecutionInfoFormatter", + "io.r2dbc.proxy.support.QueryExecutionInfoFormatter", + "io.r2dbc.proxy.util.Assert", + packageName + ".R2dbcDecorator", + packageName + ".R2dbcSqlCommentInjector", + packageName + ".R2dbcTracingSupport", + packageName + ".R2dbcTracingSupport$ConnectionMetadataListener", + packageName + ".TraceProxyExecutionListener", + }; + } + + @Override + public void methodAdvice(MethodTransformer transformer) { + transformer.applyAdvice( + isMethod() + .and(named("invoke")) + .and(takesArguments(3)) + .and(takesArgument(0, Object.class)) + .and(takesArgument(1, Method.class)) + .and(takesArgument(2, Object[].class)), + getClass().getName() + "$InvokeAdvice"); + } + + public static class InvokeAdvice { + + @Advice.OnMethodEnter(suppress = Throwable.class) + public static void onEnter( + @Advice.Argument(1) final Method method, + @Advice.Argument(value = 2, readOnly = false) Object[] args, + @Advice.FieldValue("connectionInfo") final ConnectionInfo connectionInfo) { + if (args == null || args.length == 0) { + return; + } + if (!"createStatement".equals(method.getName())) { + return; + } + if (!(args[0] instanceof String)) { + return; + } + + String dbmMode = Config.get().getDbmPropagationMode(); + boolean injectComment = + Config.DBM_PROPAGATION_MODE_FULL.equals(dbmMode) + || Config.DBM_PROPAGATION_MODE_STATIC.equals(dbmMode) + || Config.DBM_PROPAGATION_MODE_DYNAMIC_SERVICE.equals(dbmMode); + if (!injectComment) { + return; + } + + String sql = (String) args[0]; + + // Look up connection metadata from the map maintained by R2dbcTracingSupport + String hostname = null; + String dbName = null; + String dbService = null; + String dbType = null; + + ConnectionFactoryOptions options = R2dbcTracingSupport.CONNECTION_OPTIONS.get(connectionInfo); + if (options != null) { + dbType = DECORATE.extractDbType(options); + dbService = DECORATE.getDbService(options); + CharSequence hostnameSeq = null; + if (options.hasOption(ConnectionFactoryOptions.HOST)) { + Object host = options.getValue(ConnectionFactoryOptions.HOST); + if (host != null) { + hostnameSeq = host.toString(); + } + } + hostname = hostnameSeq != null ? hostnameSeq.toString() : null; + if (options.hasOption(ConnectionFactoryOptions.DATABASE)) { + Object db = options.getValue(ConnectionFactoryOptions.DATABASE); + dbName = db != null ? db.toString() : null; + } + } + + String injected = R2dbcSqlCommentInjector.inject(sql, dbService, dbType, hostname, dbName); + if (!sql.equals(injected)) { + // Replace the SQL argument with the injected version. + // We must create a new array because ByteBuddy advice cannot mutate the original + // array reference in place for @Advice.Argument(readOnly=false). + Object[] newArgs = new Object[args.length]; + System.arraycopy(args, 0, newArgs, 0, args.length); + newArgs[0] = injected; + args = newArgs; + } + } + } +} diff --git a/dd-java-agent/instrumentation/r2dbc/r2dbc-1.0/src/main/java/datadog/trace/instrumentation/r2dbc/R2dbcDecorator.java b/dd-java-agent/instrumentation/r2dbc/r2dbc-1.0/src/main/java/datadog/trace/instrumentation/r2dbc/R2dbcDecorator.java new file mode 100644 index 00000000000..dd4ea8881ea --- /dev/null +++ b/dd-java-agent/instrumentation/r2dbc/r2dbc-1.0/src/main/java/datadog/trace/instrumentation/r2dbc/R2dbcDecorator.java @@ -0,0 +1,105 @@ +package datadog.trace.instrumentation.r2dbc; + +import datadog.trace.api.naming.SpanNaming; +import datadog.trace.bootstrap.instrumentation.api.InternalSpanTypes; +import datadog.trace.bootstrap.instrumentation.api.UTF8BytesString; +import datadog.trace.bootstrap.instrumentation.decorator.DatabaseClientDecorator; +import io.r2dbc.spi.ConnectionFactoryOptions; + +public class R2dbcDecorator extends DatabaseClientDecorator { + + public static final R2dbcDecorator DECORATE = new R2dbcDecorator(); + + static final CharSequence R2DBC_QUERY = + UTF8BytesString.create(SpanNaming.instance().namingSchema().database().operation("r2dbc")); + private static final CharSequence R2DBC = UTF8BytesString.create("r2dbc"); + private static final String DEFAULT_SERVICE_NAME = + SpanNaming.instance().namingSchema().database().service("r2dbc"); + + @Override + protected String[] instrumentationNames() { + return new String[] {"r2dbc"}; + } + + @Override + protected String service() { + return DEFAULT_SERVICE_NAME; + } + + @Override + protected CharSequence component() { + return R2DBC; + } + + @Override + protected CharSequence spanType() { + return InternalSpanTypes.SQL; + } + + @Override + protected String dbType() { + return "r2dbc"; + } + + @Override + protected String dbUser(ConnectionFactoryOptions options) { + if (options == null) { + return null; + } + Object user = options.getValue(ConnectionFactoryOptions.USER); + return user != null ? user.toString() : null; + } + + @Override + protected String dbInstance(ConnectionFactoryOptions options) { + if (options == null) { + return null; + } + Object database = options.getValue(ConnectionFactoryOptions.DATABASE); + return database != null ? database.toString() : null; + } + + @Override + protected CharSequence dbHostname(ConnectionFactoryOptions options) { + if (options == null) { + return null; + } + Object host = options.getValue(ConnectionFactoryOptions.HOST); + return host != null ? host.toString() : null; + } + + public String extractDbType(ConnectionFactoryOptions options) { + if (options != null && options.hasOption(ConnectionFactoryOptions.DRIVER)) { + Object driver = options.getValue(ConnectionFactoryOptions.DRIVER); + if (driver != null) { + return driver.toString(); + } + } + return "r2dbc"; + } + + /** Exposes the protected {@link #processDatabaseType} for use by the listener. */ + public void applyDatabaseType( + datadog.trace.bootstrap.instrumentation.api.AgentSpan span, String dbType) { + processDatabaseType(span, dbType); + } + + /** + * Returns the database service name derived from the connection options. Used for DBM SQL comment + * injection. + */ + public String getDbService(ConnectionFactoryOptions options) { + String dbType = extractDbType(options); + String instanceName = dbInstance(options); + return dbService(dbType, instanceName); + } + + @Override + protected void postProcessServiceAndOperationName( + datadog.trace.bootstrap.instrumentation.api.AgentSpan span, NamingEntry namingEntry) { + if (namingEntry.getService() != null) { + span.setServiceName(namingEntry.getService(), component()); + } + span.setOperationName(namingEntry.getOperation()); + } +} diff --git a/dd-java-agent/instrumentation/r2dbc/r2dbc-1.0/src/main/java/datadog/trace/instrumentation/r2dbc/R2dbcInstrumentation.java b/dd-java-agent/instrumentation/r2dbc/r2dbc-1.0/src/main/java/datadog/trace/instrumentation/r2dbc/R2dbcInstrumentation.java new file mode 100644 index 00000000000..09a3b8a2212 --- /dev/null +++ b/dd-java-agent/instrumentation/r2dbc/r2dbc-1.0/src/main/java/datadog/trace/instrumentation/r2dbc/R2dbcInstrumentation.java @@ -0,0 +1,155 @@ +package datadog.trace.instrumentation.r2dbc; + +import static datadog.trace.agent.tooling.bytebuddy.matcher.NameMatchers.named; +import static net.bytebuddy.matcher.ElementMatchers.isMethod; +import static net.bytebuddy.matcher.ElementMatchers.isStatic; +import static net.bytebuddy.matcher.ElementMatchers.takesArgument; +import static net.bytebuddy.matcher.ElementMatchers.takesArguments; + +import com.google.auto.service.AutoService; +import datadog.trace.agent.tooling.Instrumenter; +import datadog.trace.agent.tooling.InstrumenterModule; +import io.r2dbc.spi.ConnectionFactory; +import io.r2dbc.spi.ConnectionFactoryOptions; +import net.bytebuddy.asm.Advice; + +@AutoService(InstrumenterModule.class) +public class R2dbcInstrumentation extends InstrumenterModule.Tracing + implements Instrumenter.ForSingleType, Instrumenter.HasMethodAdvice { + + public R2dbcInstrumentation() { + super("r2dbc"); + } + + @Override + public String instrumentedType() { + return "io.r2dbc.spi.ConnectionFactories"; + } + + @Override + public String[] helperClassNames() { + return new String[] { + // r2dbc-proxy is bundled (implementation, not compileOnly — see build.gradle) but its + // classes must ALSO be listed here: ordinary R2DBC apps don't depend on r2dbc-proxy, so + // their classloader can't resolve it. helperClassNames() injects these classes directly + // into the target app's classloader alongside our own, satisfying the runtime muzzle + // safety check AND avoiding NoClassDefFoundError at actual call time — r2dbc-proxy's + // internal call graph (ProxyConnectionFactory.builder() -> ProxyConfig -> callback/* + // -> core/*, util/*) reaches far more classes than the handful directly imported by + // our own helper classes, so the full non-optional class set is listed here rather + // than a hand-picked subset (io.r2dbc.proxy.observation.* is excluded: it's an + // optional Micrometer-Observation integration this module never uses, and requiring + // it would add an undeclared io.micrometer dependency). + // + // ORDER MATTERS: helpers are injected in list order via ClassLoader.defineClass(), + // which requires a type's supertypes/superinterfaces to already be defined on that + // classloader. This is a genuine topological sort of the class's extends/implements + // graph (computed from javap output, not just grouped by package/alphabetized — + // alphabetizing within a package breaks e.g. ValueStore-before-DefaultValueStore). + "io.r2dbc.proxy.callback.AfterQueryCallbackInvoker", + "io.r2dbc.proxy.callback.CallbackHandler", + "io.r2dbc.proxy.callback.CallbackHandlerSupport", + "io.r2dbc.proxy.callback.BatchCallbackHandler", + "io.r2dbc.proxy.callback.CallbackHandlerSupport$MethodInvocationStrategy", + "io.r2dbc.proxy.callback.ConnectionCallbackHandler", + "io.r2dbc.proxy.callback.ConnectionFactoryCallbackHandler", + "io.r2dbc.proxy.callback.MethodInvocationSubscriber", + "io.r2dbc.proxy.callback.ConnectionFactoryCreateMethodInvocationSubscriber", + "io.r2dbc.proxy.callback.ConnectionHolder", + "io.r2dbc.proxy.callback.ConnectionIdManager", + "io.r2dbc.proxy.callback.DefaultConnectionIdManager", + "io.r2dbc.proxy.core.ConnectionInfo", + "io.r2dbc.proxy.callback.DefaultConnectionInfo", + "io.r2dbc.proxy.callback.DelegatingContextView", + "io.r2dbc.proxy.callback.ProxyFactory", + "io.r2dbc.proxy.callback.JdkProxyFactory", + "io.r2dbc.proxy.callback.JdkProxyFactory$CallbackInvocationHandler", + "io.r2dbc.proxy.callback.ProxyFactoryFactory", + "io.r2dbc.proxy.callback.JdkProxyFactoryFactory", + "io.r2dbc.proxy.core.BindInfo", + "io.r2dbc.proxy.callback.MutableBindInfo", + "io.r2dbc.proxy.core.MethodExecutionInfo", + "io.r2dbc.proxy.callback.MutableMethodExecutionInfo", + "io.r2dbc.proxy.core.QueryExecutionInfo", + "io.r2dbc.proxy.callback.MutableQueryExecutionInfo", + "io.r2dbc.proxy.core.StatementInfo", + "io.r2dbc.proxy.callback.MutableStatementInfo", + "io.r2dbc.proxy.callback.ProxyConfig", + "io.r2dbc.proxy.callback.ProxyConfig$1", + "io.r2dbc.proxy.callback.ProxyConfig$Builder", + "io.r2dbc.proxy.callback.ProxyConfigHolder", + "io.r2dbc.proxy.callback.ProxyUtils", + "io.r2dbc.proxy.callback.QueriesExecutionContext", + "io.r2dbc.proxy.callback.QueryInvocationSubscriber", + "io.r2dbc.proxy.callback.ResultCallbackHandler", + "io.r2dbc.proxy.callback.ResultInvocationSubscriber", + "io.r2dbc.proxy.callback.RowCallbackHandler", + "io.r2dbc.proxy.callback.StatementCallbackHandler", + "io.r2dbc.proxy.callback.StopWatch", + "io.r2dbc.proxy.core.Binding", + "io.r2dbc.proxy.core.Bindings", + "io.r2dbc.proxy.core.Bindings$1", + "io.r2dbc.proxy.core.Bindings$IndexBinding", + "io.r2dbc.proxy.core.Bindings$NamedBinding", + "io.r2dbc.proxy.core.BoundValue", + "io.r2dbc.proxy.core.BoundValue$DefaultBoundValue", + "io.r2dbc.proxy.core.ValueStore", + "io.r2dbc.proxy.core.DefaultValueStore", + "io.r2dbc.proxy.core.ExecutionType", + "io.r2dbc.proxy.core.ProxyEventType", + "io.r2dbc.proxy.core.QueryInfo", + "io.r2dbc.proxy.core.R2dbcProxyException", + "io.r2dbc.proxy.listener.BindParameterConverter", + "io.r2dbc.proxy.listener.BindParameterConverter$1", + "io.r2dbc.proxy.listener.BindParameterConverter$BindOperation", + "io.r2dbc.proxy.listener.ProxyExecutionListener", + "io.r2dbc.proxy.listener.CompositeProxyExecutionListener", + "io.r2dbc.proxy.listener.LastExecutionAwareListener", + "io.r2dbc.proxy.listener.ProxyMethodExecutionListener", + "io.r2dbc.proxy.listener.ProxyMethodExecutionListenerAdapter", + "io.r2dbc.proxy.listener.ResultRowConverter", + "io.r2dbc.proxy.listener.ResultRowConverter$GetOperation", + "io.r2dbc.proxy.ProxyConnectionFactory", + "io.r2dbc.proxy.ProxyConnectionFactory$1", + "io.r2dbc.proxy.ProxyConnectionFactory$Builder", + "io.r2dbc.proxy.ProxyConnectionFactory$Builder$1", + "io.r2dbc.proxy.ProxyConnectionFactory$Builder$2", + "io.r2dbc.proxy.ProxyConnectionFactory$Builder$3", + "io.r2dbc.proxy.ProxyConnectionFactory$Builder$4", + "io.r2dbc.proxy.ProxyConnectionFactory$Builder$5", + "io.r2dbc.proxy.ProxyConnectionFactoryProvider", + "io.r2dbc.proxy.support.FormatterUtils", + "io.r2dbc.proxy.support.MethodExecutionInfoFormatter", + "io.r2dbc.proxy.support.QueryExecutionInfoFormatter", + "io.r2dbc.proxy.util.Assert", + packageName + ".R2dbcDecorator", + packageName + ".R2dbcSqlCommentInjector", + packageName + ".R2dbcTracingSupport", + packageName + ".R2dbcTracingSupport$ConnectionMetadataListener", + packageName + ".TraceProxyExecutionListener", + }; + } + + @Override + public void methodAdvice(MethodTransformer transformer) { + transformer.applyAdvice( + isMethod() + .and(isStatic()) + .and(named("find")) + .and(takesArguments(1)) + .and(takesArgument(0, named("io.r2dbc.spi.ConnectionFactoryOptions"))), + getClass().getName() + "$ConnectionFactoriesAdvice"); + } + + public static class ConnectionFactoriesAdvice { + + @Advice.OnMethodExit(suppress = Throwable.class) + public static void onExit( + @Advice.Return(readOnly = false) ConnectionFactory factory, + @Advice.Argument(0) ConnectionFactoryOptions options) { + if (factory != null) { + factory = R2dbcTracingSupport.wrapConnectionFactory(factory, options); + } + } + } +} diff --git a/dd-java-agent/instrumentation/r2dbc/r2dbc-1.0/src/main/java/datadog/trace/instrumentation/r2dbc/R2dbcSqlCommentInjector.java b/dd-java-agent/instrumentation/r2dbc/r2dbc-1.0/src/main/java/datadog/trace/instrumentation/r2dbc/R2dbcSqlCommentInjector.java new file mode 100644 index 00000000000..20374e7e7ad --- /dev/null +++ b/dd-java-agent/instrumentation/r2dbc/r2dbc-1.0/src/main/java/datadog/trace/instrumentation/r2dbc/R2dbcSqlCommentInjector.java @@ -0,0 +1,82 @@ +package datadog.trace.instrumentation.r2dbc; + +import static datadog.trace.bootstrap.instrumentation.api.AgentTracer.activeSpan; + +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; + +/** + * Injects DBM SQL comments into R2DBC queries. Reuses {@link SharedDBCommenter} to build the + * comment content (service metadata, trace context) and wraps it in SQL comment delimiters. + * + *

This is the R2DBC equivalent of JDBC's {@code SQLCommenter}. It is intentionally simpler + * because R2DBC does not have the same edge cases (callable statements, pg_hint_plan) as JDBC. + */ +public final class R2dbcSqlCommentInjector { + + private static final String OPEN_COMMENT = "/*"; + private static final String CLOSE_COMMENT = "*/"; + + private R2dbcSqlCommentInjector() {} + + /** + * Injects a DBM SQL comment into the given query string if DBM propagation is enabled. + * + * @param sql the original SQL query + * @param dbService the database service name for dddbs + * @param dbType the database type (e.g. "h2", "postgresql") + * @param hostname the database hostname + * @param dbName the database name + * @return the SQL with injected comment, or the original SQL if DBM is disabled + */ + public static String inject( + String sql, String dbService, String dbType, String hostname, String dbName) { + if (sql == null || sql.isEmpty()) { + return sql; + } + + String dbmMode = Config.get().getDbmPropagationMode(); + boolean injectComment = + Config.DBM_PROPAGATION_MODE_FULL.equals(dbmMode) + || Config.DBM_PROPAGATION_MODE_STATIC.equals(dbmMode) + || Config.DBM_PROPAGATION_MODE_DYNAMIC_SERVICE.equals(dbmMode); + + if (!injectComment) { + return sql; + } + + // Generate traceparent only in full mode + String traceParent = null; + if (Config.DBM_PROPAGATION_MODE_FULL.equals(dbmMode)) { + AgentSpan activeSpan = activeSpan(); + if (activeSpan != null) { + Integer priority = activeSpan.forceSamplingDecision(); + if (priority != null) { + traceParent = W3CTraceParent.from(activeSpan); + } + } + } + + String commentContent = + SharedDBCommenter.buildComment(dbService, dbType, hostname, dbName, traceParent); + if (commentContent == null) { + return sql; + } + + // Check for existing DD comment to avoid duplicate injection + if (sql.startsWith(OPEN_COMMENT) && SharedDBCommenter.containsTraceComment(sql)) { + return sql; + } + + // Prepend the comment to the SQL query + StringBuilder sb = new StringBuilder(sql.length() + commentContent.length() + 6); + sb.append(OPEN_COMMENT); + sb.append(commentContent); + sb.append(CLOSE_COMMENT); + sb.append(' '); + sb.append(sql); + return sb.toString(); + } +} diff --git a/dd-java-agent/instrumentation/r2dbc/r2dbc-1.0/src/main/java/datadog/trace/instrumentation/r2dbc/R2dbcTracingSupport.java b/dd-java-agent/instrumentation/r2dbc/r2dbc-1.0/src/main/java/datadog/trace/instrumentation/r2dbc/R2dbcTracingSupport.java new file mode 100644 index 00000000000..1fe1d48f6a0 --- /dev/null +++ b/dd-java-agent/instrumentation/r2dbc/r2dbc-1.0/src/main/java/datadog/trace/instrumentation/r2dbc/R2dbcTracingSupport.java @@ -0,0 +1,71 @@ +package datadog.trace.instrumentation.r2dbc; + +import io.r2dbc.proxy.ProxyConnectionFactory; +import io.r2dbc.proxy.core.ConnectionInfo; +import io.r2dbc.proxy.listener.ProxyExecutionListener; +import io.r2dbc.spi.ConnectionFactory; +import io.r2dbc.spi.ConnectionFactoryOptions; +import java.util.Map; +import java.util.concurrent.ConcurrentHashMap; + +/** + * Wraps a {@link ConnectionFactory} with r2dbc-proxy to install a tracing listener. Also maintains + * a mapping from {@link ConnectionInfo} to {@link ConnectionFactoryOptions} so that DBM SQL comment + * injection can access connection metadata (host, database, driver type). + */ +public final class R2dbcTracingSupport { + + /** + * Maps R2DBC proxy {@link ConnectionInfo} instances to the {@link ConnectionFactoryOptions} used + * to create the connection factory. This allows the DBM SQL comment injector to resolve + * connection metadata (hostname, database name, db type) when intercepting {@code + * createStatement} calls. + * + *

Entries are added when the proxy listener's {@code afterMethod} fires for {@code create()} + * (connection creation), and removed when the connection is closed. + */ + static final Map CONNECTION_OPTIONS = + new ConcurrentHashMap<>(); + + private R2dbcTracingSupport() {} + + public static ConnectionFactory wrapConnectionFactory( + ConnectionFactory factory, ConnectionFactoryOptions options) { + TraceProxyExecutionListener queryListener = new TraceProxyExecutionListener(options); + ConnectionMetadataListener metadataListener = new ConnectionMetadataListener(options); + + return ProxyConnectionFactory.builder(factory) + .listener(queryListener) + .listener(metadataListener) + .build(); + } + + /** + * A lightweight listener that tracks connection creation and close events to maintain the {@link + * #CONNECTION_OPTIONS} map. This allows {@link R2dbcConnectionCallbackInstrumentation} to look up + * connection metadata when injecting SQL comments. + */ + static final class ConnectionMetadataListener implements ProxyExecutionListener { + private final ConnectionFactoryOptions options; + + ConnectionMetadataListener(ConnectionFactoryOptions options) { + this.options = options; + } + + @Override + public void afterMethod(io.r2dbc.proxy.core.MethodExecutionInfo execInfo) { + String methodName = execInfo.getMethod().getName(); + ConnectionInfo connInfo = execInfo.getConnectionInfo(); + if (connInfo == null) { + return; + } + if ("create".equals(methodName) && execInfo.getThrown() == null) { + // Connection was successfully created — register the metadata + CONNECTION_OPTIONS.put(connInfo, options); + } else if ("close".equals(methodName)) { + // Connection closed — clean up + CONNECTION_OPTIONS.remove(connInfo); + } + } + } +} diff --git a/dd-java-agent/instrumentation/r2dbc/r2dbc-1.0/src/main/java/datadog/trace/instrumentation/r2dbc/TraceProxyExecutionListener.java b/dd-java-agent/instrumentation/r2dbc/r2dbc-1.0/src/main/java/datadog/trace/instrumentation/r2dbc/TraceProxyExecutionListener.java new file mode 100644 index 00000000000..ececd321e04 --- /dev/null +++ b/dd-java-agent/instrumentation/r2dbc/r2dbc-1.0/src/main/java/datadog/trace/instrumentation/r2dbc/TraceProxyExecutionListener.java @@ -0,0 +1,93 @@ +package datadog.trace.instrumentation.r2dbc; + +import static datadog.trace.bootstrap.instrumentation.api.AgentTracer.startSpan; +import static datadog.trace.instrumentation.r2dbc.R2dbcDecorator.DECORATE; +import static datadog.trace.instrumentation.r2dbc.R2dbcDecorator.R2DBC_QUERY; + +import datadog.trace.api.Config; +import datadog.trace.bootstrap.instrumentation.api.AgentSpan; +import io.r2dbc.proxy.core.QueryExecutionInfo; +import io.r2dbc.proxy.core.QueryInfo; +import io.r2dbc.proxy.listener.ProxyExecutionListener; +import io.r2dbc.spi.ConnectionFactoryOptions; +import java.util.List; + +/** + * R2DBC proxy listener that creates database spans around query executions. The r2dbc-proxy + * framework owns the reactive lifecycle (complete/error/cancel), so this listener does not need to + * handle cancellation — the {@code afterQuery} callback fires in all cases. + * + *

When Database Monitoring (DBM) is enabled via {@code dd.dbm.propagation.mode}, this listener + * also sets the {@code _dd.dbm_trace_injected} tag on spans. SQL comment injection is handled + * separately by {@link R2dbcConnectionCallbackInstrumentation}. + */ +public final class TraceProxyExecutionListener implements ProxyExecutionListener { + + private static final String SPAN_KEY = "datadog.span"; + private static final String DBM_TRACE_INJECTED = "_dd.dbm_trace_injected"; + + private final ConnectionFactoryOptions options; + private final boolean injectTraceContext; + + public TraceProxyExecutionListener(ConnectionFactoryOptions options) { + this.options = options; + String dbmMode = Config.get().getDbmPropagationMode(); + this.injectTraceContext = Config.DBM_PROPAGATION_MODE_FULL.equals(dbmMode); + } + + @Override + public void beforeQuery(QueryExecutionInfo execInfo) { + AgentSpan span = startSpan("r2dbc", R2DBC_QUERY); + DECORATE.afterStart(span); + + String dbType = DECORATE.extractDbType(options); + DECORATE.applyDatabaseType(span, dbType); + DECORATE.onConnection(span, options); + + String queryString = extractQuery(execInfo); + if (queryString != null) { + DECORATE.onStatement(span, queryString); + } + + if (injectTraceContext) { + Integer priority = span.forceSamplingDecision(); + if (priority != null) { + span.setTag(DBM_TRACE_INJECTED, true); + } + } + + span.setMeasured(true); + execInfo.getValueStore().put(SPAN_KEY, span); + } + + @Override + public void afterQuery(QueryExecutionInfo execInfo) { + AgentSpan span = execInfo.getValueStore().get(SPAN_KEY, AgentSpan.class); + if (span == null) { + return; + } + if (execInfo.getThrowable() != null) { + DECORATE.onError(span, execInfo.getThrowable()); + } + DECORATE.beforeFinish(span); + span.finish(); + } + + private static String extractQuery(QueryExecutionInfo execInfo) { + List queries = execInfo.getQueries(); + if (queries == null || queries.isEmpty()) { + return null; + } + if (queries.size() == 1) { + return queries.get(0).getQuery(); + } + StringBuilder sb = new StringBuilder(); + for (int i = 0; i < queries.size(); i++) { + if (i > 0) { + sb.append("; "); + } + sb.append(queries.get(i).getQuery()); + } + return sb.toString(); + } +} diff --git a/dd-java-agent/instrumentation/r2dbc/r2dbc-1.0/src/test/java/datadog/trace/instrumentation/r2dbc/R2dbcDbmForkedTest.java b/dd-java-agent/instrumentation/r2dbc/r2dbc-1.0/src/test/java/datadog/trace/instrumentation/r2dbc/R2dbcDbmForkedTest.java new file mode 100644 index 00000000000..512d7063a17 --- /dev/null +++ b/dd-java-agent/instrumentation/r2dbc/r2dbc-1.0/src/test/java/datadog/trace/instrumentation/r2dbc/R2dbcDbmForkedTest.java @@ -0,0 +1,276 @@ +package datadog.trace.instrumentation.r2dbc; + +import static datadog.trace.agent.test.assertions.SpanMatcher.span; +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.SORT_BY_START_TIME; +import static datadog.trace.agent.test.assertions.TraceMatcher.trace; +import static datadog.trace.bootstrap.instrumentation.api.AgentTracer.activateSpan; +import static datadog.trace.bootstrap.instrumentation.api.AgentTracer.startSpan; +import static datadog.trace.instrumentation.r2dbc.R2dbcInstrumentationTest.eqs; +import static datadog.trace.test.junit.utils.assertions.Matchers.any; +import static datadog.trace.test.junit.utils.assertions.Matchers.is; + +import datadog.trace.agent.test.AbstractInstrumentationTest; +import datadog.trace.api.DDSpanTypes; +import datadog.trace.api.DDTags; +import datadog.trace.bootstrap.instrumentation.api.AgentScope; +import datadog.trace.bootstrap.instrumentation.api.AgentSpan; +import datadog.trace.bootstrap.instrumentation.api.Tags; +import datadog.trace.test.junit.utils.config.WithConfig; +import io.r2dbc.spi.Connection; +import io.r2dbc.spi.ConnectionFactories; +import io.r2dbc.spi.ConnectionFactory; +import io.r2dbc.spi.ConnectionFactoryOptions; +import io.r2dbc.spi.Result; +import java.util.regex.Pattern; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import reactor.core.publisher.Flux; +import reactor.core.publisher.Mono; + +/** + * Tests for R2DBC Database Monitoring (DBM) feature. Verifies that connection metadata tags + * (db.instance, db.user, peer.hostname) are correctly populated on spans, and that the + * _dd.dbm_trace_injected tag is set when DBM propagation mode is "full". + */ +@WithConfig(key = "dbm.propagation.mode", value = "full") +@WithConfig(key = "service", value = "test_service", addPrefix = false) +class R2dbcDbmForkedTest extends AbstractInstrumentationTest { + + private static final Pattern H2_QUERY = Pattern.compile("h2\\.query"); + + private ConnectionFactory connectionFactory; + private Connection connection; + + @BeforeEach + public void setUp() { + connectionFactory = + ConnectionFactories.get( + ConnectionFactoryOptions.builder() + .option(ConnectionFactoryOptions.DRIVER, "h2") + .option(ConnectionFactoryOptions.PROTOCOL, "mem") + .option(ConnectionFactoryOptions.DATABASE, "testdb") + .build()); + connection = Mono.from(connectionFactory.create()).block(); + + // Create a test table + Mono.from( + connection + .createStatement( + "CREATE TABLE IF NOT EXISTS test_table (id INT, name VARCHAR(255))") + .execute()) + .flatMapMany(result -> result.getRowsUpdated()) + .blockLast(); + + // Clear any traces from setup + tracer.flush(); + writer.clear(); + } + + @AfterEach + public void tearDown() { + if (connection != null) { + Mono.from(connection.createStatement("DROP TABLE IF EXISTS test_table").execute()) + .flatMapMany(result -> result.getRowsUpdated()) + .blockLast(); + Mono.from(connection.close()).block(); + } + } + + @Test + void dbmPopulatesConnectionMetadataTagsOnSelectQuery() { + AgentSpan parent = startSpan("test", "parent"); + try (AgentScope scope = activateSpan(parent)) { + Flux.from(connection.createStatement("SELECT * FROM test_table").execute()) + .flatMap(result -> result.map((row, metadata) -> row.get(0))) + .collectList() + .block(); + } finally { + parent.finish(); + } + + assertTraces( + trace( + SORT_BY_START_TIME, + span().root().operationName("parent"), + span() + .childOfPrevious() + .operationName(H2_QUERY) + .resourceName("SELECT * FROM test_table") + .type(DDSpanTypes.SQL) + .measured() + .tags( + tag(Tags.COMPONENT, eqs("r2dbc")), + tag(Tags.SPAN_KIND, eqs(Tags.SPAN_KIND_CLIENT)), + tag(Tags.DB_TYPE, eqs("h2")), + tag(Tags.DB_INSTANCE, eqs("testdb")), + tag("_dd.dbm_trace_injected", is(true)), + tag("_dd.svc_src", any()), + defaultTags()))); + } + + @Test + void dbmSetsTraceInjectedTagInFullMode() { + AgentSpan parent = startSpan("test", "parent"); + try (AgentScope scope = activateSpan(parent)) { + Flux.from(connection.createStatement("SELECT * FROM test_table").execute()) + .flatMap(result -> result.map((row, metadata) -> row.get(0))) + .collectList() + .block(); + } finally { + parent.finish(); + } + + // In full mode, the span should have the _dd.dbm_trace_injected tag set to true + assertTraces( + trace( + SORT_BY_START_TIME, + span().root().operationName("parent"), + span() + .childOfPrevious() + .operationName(H2_QUERY) + .resourceName("SELECT * FROM test_table") + .type(DDSpanTypes.SQL) + .measured() + .tags( + tag(Tags.COMPONENT, eqs("r2dbc")), + tag(Tags.SPAN_KIND, eqs(Tags.SPAN_KIND_CLIENT)), + tag(Tags.DB_TYPE, eqs("h2")), + tag(Tags.DB_INSTANCE, eqs("testdb")), + tag("_dd.dbm_trace_injected", is(true)), + tag("_dd.svc_src", any()), + defaultTags()))); + } + + @Test + void dbmPopulatesTagsOnInsertQuery() { + AgentSpan parent = startSpan("test", "parent"); + try (AgentScope scope = activateSpan(parent)) { + Mono.from( + connection + .createStatement("INSERT INTO test_table (id, name) VALUES (1, 'test')") + .execute()) + .flatMapMany(Result::getRowsUpdated) + .blockLast(); + } finally { + parent.finish(); + } + + assertTraces( + trace( + SORT_BY_START_TIME, + span().root().operationName("parent"), + span() + .childOfPrevious() + .operationName(H2_QUERY) + .resourceName(Pattern.compile("INSERT INTO test_table.*")) + .type(DDSpanTypes.SQL) + .measured() + .tags( + tag(Tags.COMPONENT, eqs("r2dbc")), + tag(Tags.SPAN_KIND, eqs(Tags.SPAN_KIND_CLIENT)), + tag(Tags.DB_TYPE, eqs("h2")), + tag(Tags.DB_INSTANCE, eqs("testdb")), + tag("_dd.dbm_trace_injected", is(true)), + tag("_dd.svc_src", any()), + defaultTags()))); + } + + @Test + void dbmPreservesErrorTagsOnFailedQuery() { + AgentSpan parent = startSpan("test", "parent"); + try (AgentScope scope = activateSpan(parent)) { + try { + Flux.from(connection.createStatement("SELECT * FROM nonexistent_table").execute()) + .flatMap(result -> result.map((row, metadata) -> row.get(0))) + .collectList() + .block(); + } catch (Exception ignored) { + // Expected to fail + } + } finally { + parent.finish(); + } + + // Error spans should still work correctly with DBM enabled, and still carry DBM tags + assertTraces( + trace( + SORT_BY_START_TIME, + span().root().operationName("parent"), + span() + .childOfPrevious() + .operationName(H2_QUERY) + .resourceName("SELECT * FROM nonexistent_table") + .type(DDSpanTypes.SQL) + .error() + .measured() + .tags( + tag(Tags.COMPONENT, eqs("r2dbc")), + tag(Tags.SPAN_KIND, eqs(Tags.SPAN_KIND_CLIENT)), + tag(Tags.DB_TYPE, eqs("h2")), + tag(Tags.DB_INSTANCE, eqs("testdb")), + tag("_dd.dbm_trace_injected", is(true)), + tag("_dd.svc_src", any()), + tag(DDTags.ERROR_MSG, any()), + error(Exception.class), + defaultTags()))); + } + + @Test + void dbmWorksAcrossMultipleQueries() { + AgentSpan parent = startSpan("test", "parent"); + try (AgentScope scope = activateSpan(parent)) { + // First query: INSERT + Mono.from( + connection + .createStatement("INSERT INTO test_table (id, name) VALUES (1, 'first')") + .execute()) + .flatMapMany(Result::getRowsUpdated) + .blockLast(); + + // Second query: SELECT + Flux.from(connection.createStatement("SELECT * FROM test_table").execute()) + .flatMap(result -> result.map((row, metadata) -> row.get(0))) + .collectList() + .block(); + } finally { + parent.finish(); + } + + assertTraces( + trace( + SORT_BY_START_TIME, + span().root().operationName("parent"), + span() + .childOfPrevious() + .operationName(H2_QUERY) + .resourceName(Pattern.compile("INSERT INTO test_table.*")) + .type(DDSpanTypes.SQL) + .measured() + .tags( + tag(Tags.COMPONENT, eqs("r2dbc")), + tag(Tags.SPAN_KIND, eqs(Tags.SPAN_KIND_CLIENT)), + tag(Tags.DB_TYPE, eqs("h2")), + tag(Tags.DB_INSTANCE, eqs("testdb")), + tag("_dd.dbm_trace_injected", is(true)), + tag("_dd.svc_src", any()), + defaultTags()), + span() + .childOfIndex(0) + .operationName(H2_QUERY) + .resourceName("SELECT * FROM test_table") + .type(DDSpanTypes.SQL) + .measured() + .tags( + tag(Tags.COMPONENT, eqs("r2dbc")), + tag(Tags.SPAN_KIND, eqs(Tags.SPAN_KIND_CLIENT)), + tag(Tags.DB_TYPE, eqs("h2")), + tag(Tags.DB_INSTANCE, eqs("testdb")), + tag("_dd.dbm_trace_injected", is(true)), + tag("_dd.svc_src", any()), + defaultTags()))); + } +} diff --git a/dd-java-agent/instrumentation/r2dbc/r2dbc-1.0/src/test/java/datadog/trace/instrumentation/r2dbc/R2dbcInstrumentationTest.java b/dd-java-agent/instrumentation/r2dbc/r2dbc-1.0/src/test/java/datadog/trace/instrumentation/r2dbc/R2dbcInstrumentationTest.java new file mode 100644 index 00000000000..69687904145 --- /dev/null +++ b/dd-java-agent/instrumentation/r2dbc/r2dbc-1.0/src/test/java/datadog/trace/instrumentation/r2dbc/R2dbcInstrumentationTest.java @@ -0,0 +1,329 @@ +package datadog.trace.instrumentation.r2dbc; + +import static datadog.trace.agent.test.assertions.SpanMatcher.span; +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.SORT_BY_START_TIME; +import static datadog.trace.agent.test.assertions.TraceMatcher.trace; +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 datadog.trace.agent.test.AbstractInstrumentationTest; +import datadog.trace.api.DDSpanTypes; +import datadog.trace.api.DDTags; +import datadog.trace.bootstrap.instrumentation.api.AgentScope; +import datadog.trace.bootstrap.instrumentation.api.AgentSpan; +import datadog.trace.bootstrap.instrumentation.api.Tags; +import datadog.trace.test.junit.utils.assertions.Matcher; +import io.r2dbc.spi.Connection; +import io.r2dbc.spi.ConnectionFactories; +import io.r2dbc.spi.ConnectionFactory; +import io.r2dbc.spi.Result; +import java.util.Optional; +import java.util.regex.Pattern; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import reactor.core.publisher.Flux; +import reactor.core.publisher.Mono; + +/** + * Instrumentation tests for the R2DBC listener-based integration. Verifies that the tracer creates + * database spans when queries are executed through R2DBC's reactive connection API. + */ +class R2dbcInstrumentationTest extends AbstractInstrumentationTest { + + private static final Pattern H2_QUERY = Pattern.compile("h2\\.query"); + + /** + * Creates a matcher that compares by {@code toString()} to handle both {@code String} and {@code + * UTF8BytesString} values in tag comparisons. + */ + @SuppressWarnings("unchecked") + static Matcher eqs(String expected) { + return new Matcher() { + @Override + public Optional expected() { + return Optional.of((T) expected); + } + + @Override + public String failureReason() { + return "Unexpected value"; + } + + @Override + public boolean test(T t) { + return t != null && expected.equals(t.toString()); + } + }; + } + + private ConnectionFactory connectionFactory; + private Connection connection; + + @BeforeEach + public void setUp() { + connectionFactory = ConnectionFactories.get("r2dbc:h2:mem:///testdb;DB_CLOSE_DELAY=-1"); + connection = Mono.from(connectionFactory.create()).block(); + + // Create a test table + Mono.from( + connection + .createStatement( + "CREATE TABLE IF NOT EXISTS test_table (id INT, name VARCHAR(255))") + .execute()) + .flatMapMany(result -> result.getRowsUpdated()) + .blockLast(); + + // Clear any traces from setup + tracer.flush(); + writer.clear(); + } + + @AfterEach + public void tearDown() { + if (connection != null) { + Mono.from(connection.createStatement("DROP TABLE IF EXISTS test_table").execute()) + .flatMapMany(result -> result.getRowsUpdated()) + .blockLast(); + Mono.from(connection.close()).block(); + } + } + + @Test + void selectQueryCreatesSpan() { + AgentSpan parent = startSpan("test", "parent"); + try (AgentScope scope = activateSpan(parent)) { + Flux.from(connection.createStatement("SELECT * FROM test_table").execute()) + .flatMap(result -> result.map((row, metadata) -> row.get(0))) + .collectList() + .block(); + } finally { + parent.finish(); + } + + assertTraces( + trace( + SORT_BY_START_TIME, + span().root().operationName("parent"), + span() + .childOfPrevious() + .operationName(H2_QUERY) + .resourceName("SELECT * FROM test_table") + .type(DDSpanTypes.SQL) + .measured() + .tags( + tag(Tags.COMPONENT, eqs("r2dbc")), + tag(Tags.SPAN_KIND, eqs(Tags.SPAN_KIND_CLIENT)), + tag(Tags.DB_TYPE, eqs("h2")), + tag(Tags.DB_INSTANCE, any()), + tag("_dd.svc_src", any()), + defaultTags()))); + } + + @Test + void insertQueryCreatesSpan() { + AgentSpan parent = startSpan("test", "parent"); + try (AgentScope scope = activateSpan(parent)) { + Mono.from( + connection + .createStatement("INSERT INTO test_table (id, name) VALUES (1, 'test')") + .execute()) + .flatMapMany(Result::getRowsUpdated) + .blockLast(); + } finally { + parent.finish(); + } + + assertTraces( + trace( + SORT_BY_START_TIME, + span().root().operationName("parent"), + span() + .childOfPrevious() + .operationName(H2_QUERY) + .resourceName(Pattern.compile("INSERT INTO test_table.*")) + .type(DDSpanTypes.SQL) + .measured() + .tags( + tag(Tags.COMPONENT, eqs("r2dbc")), + tag(Tags.SPAN_KIND, eqs(Tags.SPAN_KIND_CLIENT)), + tag(Tags.DB_TYPE, eqs("h2")), + tag(Tags.DB_INSTANCE, any()), + tag("_dd.svc_src", any()), + defaultTags()))); + } + + @Test + void multipleQueriesCreateMultipleSpans() { + AgentSpan parent = startSpan("test", "parent"); + try (AgentScope scope = activateSpan(parent)) { + Mono.from( + connection + .createStatement("INSERT INTO test_table (id, name) VALUES (1, 'first')") + .execute()) + .flatMapMany(Result::getRowsUpdated) + .blockLast(); + + Flux.from(connection.createStatement("SELECT * FROM test_table").execute()) + .flatMap(result -> result.map((row, metadata) -> row.get(0))) + .collectList() + .block(); + } finally { + parent.finish(); + } + + assertTraces( + trace( + SORT_BY_START_TIME, + span().root().operationName("parent"), + span() + .childOfPrevious() + .operationName(H2_QUERY) + .resourceName(Pattern.compile("INSERT INTO test_table.*")) + .type(DDSpanTypes.SQL) + .measured() + .tags( + tag(Tags.COMPONENT, eqs("r2dbc")), + tag(Tags.SPAN_KIND, eqs(Tags.SPAN_KIND_CLIENT)), + tag(Tags.DB_TYPE, eqs("h2")), + tag(Tags.DB_INSTANCE, any()), + tag("_dd.svc_src", any()), + defaultTags()), + span() + .childOfIndex(0) + .operationName(H2_QUERY) + .resourceName("SELECT * FROM test_table") + .type(DDSpanTypes.SQL) + .measured() + .tags( + tag(Tags.COMPONENT, eqs("r2dbc")), + tag(Tags.SPAN_KIND, eqs(Tags.SPAN_KIND_CLIENT)), + tag(Tags.DB_TYPE, eqs("h2")), + tag(Tags.DB_INSTANCE, any()), + tag("_dd.svc_src", any()), + defaultTags()))); + } + + @Test + void errorQuerySetsErrorTags() { + AgentSpan parent = startSpan("test", "parent"); + try (AgentScope scope = activateSpan(parent)) { + try { + Flux.from(connection.createStatement("SELECT * FROM nonexistent_table").execute()) + .flatMap(result -> result.map((row, metadata) -> row.get(0))) + .collectList() + .block(); + } catch (Exception ignored) { + // Expected to fail + } + } finally { + parent.finish(); + } + + assertTraces( + trace( + SORT_BY_START_TIME, + span().root().operationName("parent"), + span() + .childOfPrevious() + .operationName(H2_QUERY) + .resourceName("SELECT * FROM nonexistent_table") + .type(DDSpanTypes.SQL) + .error() + .measured() + .tags( + tag(Tags.COMPONENT, eqs("r2dbc")), + tag(Tags.SPAN_KIND, eqs(Tags.SPAN_KIND_CLIENT)), + tag(Tags.DB_TYPE, eqs("h2")), + tag(Tags.DB_INSTANCE, any()), + tag("_dd.svc_src", any()), + tag(DDTags.ERROR_MSG, any()), + error(Exception.class), + defaultTags()))); + } + + @Test + void cancelledQueryStillFinishesSpan() { + // Insert some data first so the query has something to stream — use a separate + // span so the INSERT trace doesn't merge with the test's assertion target. + AgentSpan setupParent = startSpan("test", "setup"); + try (AgentScope setupScope = activateSpan(setupParent)) { + Mono.from( + connection + .createStatement("INSERT INTO test_table (id, name) VALUES (1, 'a')") + .execute()) + .flatMapMany(Result::getRowsUpdated) + .blockLast(); + } finally { + setupParent.finish(); + } + // Clear setup traces + tracer.flush(); + writer.clear(); + + // Now cancel a query mid-stream using take(1) + AgentSpan parent = startSpan("test", "parent"); + try (AgentScope scope = activateSpan(parent)) { + Flux.from(connection.createStatement("SELECT * FROM test_table").execute()) + .flatMap(result -> result.map((row, metadata) -> row.get(0))) + .take(1) + .blockLast(); + } finally { + parent.finish(); + } + + // The key assertion: even though the reactive stream was cancelled via take(1), + // the span must still finish — no leaked, never-finished spans. + assertTraces( + trace( + SORT_BY_START_TIME, + span().root().operationName("parent"), + span() + .childOfPrevious() + .operationName(H2_QUERY) + .resourceName("SELECT * FROM test_table") + .type(DDSpanTypes.SQL) + .measured() + .tags( + tag(Tags.COMPONENT, eqs("r2dbc")), + tag(Tags.SPAN_KIND, eqs(Tags.SPAN_KIND_CLIENT)), + tag(Tags.DB_TYPE, eqs("h2")), + tag(Tags.DB_INSTANCE, any()), + tag("_dd.svc_src", any()), + defaultTags()))); + } + + @Test + void queryWithNoActiveTraceDoesNotCreateOrphanSpans() { + // Execute a query without any active trace context + Flux.from(connection.createStatement("SELECT * FROM test_table").execute()) + .flatMap(result -> result.map((row, metadata) -> row.get(0))) + .collectList() + .block(); + + tracer.flush(); + + // The listener should still create a span (it wraps the query regardless), + // but it should be a root span rather than an orphan child. + blockUntilTracesMatch(traces -> traces.size() >= 1); + assertTraces( + trace( + span() + .root() + .operationName(H2_QUERY) + .resourceName("SELECT * FROM test_table") + .type(DDSpanTypes.SQL) + .measured() + .tags( + tag(Tags.COMPONENT, eqs("r2dbc")), + tag(Tags.SPAN_KIND, eqs(Tags.SPAN_KIND_CLIENT)), + tag(Tags.DB_TYPE, eqs("h2")), + tag(Tags.DB_INSTANCE, any()), + tag("_dd.svc_src", any()), + defaultTags()))); + } +} diff --git a/dd-java-agent/instrumentation/r2dbc/r2dbc-1.0/src/test/java/datadog/trace/instrumentation/r2dbc/R2dbcPeerServiceTest.java b/dd-java-agent/instrumentation/r2dbc/r2dbc-1.0/src/test/java/datadog/trace/instrumentation/r2dbc/R2dbcPeerServiceTest.java new file mode 100644 index 00000000000..f9206b44d19 --- /dev/null +++ b/dd-java-agent/instrumentation/r2dbc/r2dbc-1.0/src/test/java/datadog/trace/instrumentation/r2dbc/R2dbcPeerServiceTest.java @@ -0,0 +1,254 @@ +package datadog.trace.instrumentation.r2dbc; + +import static datadog.trace.agent.test.assertions.SpanMatcher.span; +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.SORT_BY_START_TIME; +import static datadog.trace.agent.test.assertions.TraceMatcher.trace; +import static datadog.trace.bootstrap.instrumentation.api.AgentTracer.activateSpan; +import static datadog.trace.bootstrap.instrumentation.api.AgentTracer.startSpan; +import static datadog.trace.instrumentation.r2dbc.R2dbcInstrumentationTest.eqs; +import static datadog.trace.test.junit.utils.assertions.Matchers.any; + +import datadog.trace.agent.test.AbstractInstrumentationTest; +import datadog.trace.api.DDSpanTypes; +import datadog.trace.api.DDTags; +import datadog.trace.bootstrap.instrumentation.api.AgentScope; +import datadog.trace.bootstrap.instrumentation.api.AgentSpan; +import datadog.trace.bootstrap.instrumentation.api.Tags; +import io.r2dbc.spi.Connection; +import io.r2dbc.spi.ConnectionFactories; +import io.r2dbc.spi.ConnectionFactory; +import io.r2dbc.spi.ConnectionFactoryOptions; +import io.r2dbc.spi.Result; +import java.util.regex.Pattern; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import reactor.core.publisher.Flux; +import reactor.core.publisher.Mono; + +/** + * Tests for R2DBC peer service feature. Verifies that the input tags used by {@code + * PeerServiceCalculator} (peer.hostname, db.instance) are correctly set on database spans when + * connection options include a HOST. Also verifies that db.user is set when USER is provided. + * + *

We use H2's in-memory mode with explicit HOST/DATABASE/USER options to exercise the full + * metadata extraction path. H2 ignores the HOST option for mem connections, but the R2DBC SPI + * stores all options so our decorator can read them and set the corresponding span tags. The + * PeerServiceCalculator will compute peer.service from these input tags — we assert on the inputs, + * not the computed output, per the peer_service feature guide. + */ +class R2dbcPeerServiceTest extends AbstractInstrumentationTest { + + private static final Pattern H2_QUERY = Pattern.compile("h2\\.query"); + + private ConnectionFactory connectionFactory; + private Connection connection; + + @BeforeEach + public void setUp() { + // Build ConnectionFactoryOptions with explicit HOST, DATABASE, and USER. + // H2's in-memory protocol ignores HOST but the R2DBC SPI stores all options, + // so our decorator can read them and set the corresponding span tags. + connectionFactory = + ConnectionFactories.get( + ConnectionFactoryOptions.builder() + .option(ConnectionFactoryOptions.DRIVER, "h2") + .option(ConnectionFactoryOptions.PROTOCOL, "mem") + .option(ConnectionFactoryOptions.HOST, "db.example.com") + .option(ConnectionFactoryOptions.DATABASE, "peerdb") + .option(ConnectionFactoryOptions.USER, "testuser") + .build()); + connection = Mono.from(connectionFactory.create()).block(); + + // Create a test table + Mono.from( + connection + .createStatement("CREATE TABLE IF NOT EXISTS peer_test (id INT, name VARCHAR(255))") + .execute()) + .flatMapMany(result -> result.getRowsUpdated()) + .blockLast(); + + // Clear any traces from setup + tracer.flush(); + writer.clear(); + } + + @AfterEach + public void tearDown() { + if (connection != null) { + Mono.from(connection.createStatement("DROP TABLE IF EXISTS peer_test").execute()) + .flatMapMany(result -> result.getRowsUpdated()) + .blockLast(); + Mono.from(connection.close()).block(); + } + } + + @Test + void peerHostnameSetOnSelectQuery() { + AgentSpan parent = startSpan("test", "parent"); + try (AgentScope scope = activateSpan(parent)) { + Flux.from(connection.createStatement("SELECT * FROM peer_test").execute()) + .flatMap(result -> result.map((row, metadata) -> row.get(0))) + .collectList() + .block(); + } finally { + parent.finish(); + } + + assertTraces( + trace( + SORT_BY_START_TIME, + span().root().operationName("parent"), + span() + .childOfPrevious() + .operationName(H2_QUERY) + .resourceName("SELECT * FROM peer_test") + .type(DDSpanTypes.SQL) + .measured() + .tags( + tag(Tags.COMPONENT, eqs("r2dbc")), + tag(Tags.SPAN_KIND, eqs(Tags.SPAN_KIND_CLIENT)), + tag(Tags.DB_TYPE, eqs("h2")), + tag(Tags.DB_INSTANCE, eqs("peerdb")), + tag(Tags.PEER_HOSTNAME, eqs("db.example.com")), + tag(Tags.DB_USER, eqs("testuser")), + tag("_dd.svc_src", any()), + defaultTags()))); + } + + @Test + void peerHostnameSetOnInsertQuery() { + AgentSpan parent = startSpan("test", "parent"); + try (AgentScope scope = activateSpan(parent)) { + Mono.from( + connection + .createStatement("INSERT INTO peer_test (id, name) VALUES (1, 'test')") + .execute()) + .flatMapMany(Result::getRowsUpdated) + .blockLast(); + } finally { + parent.finish(); + } + + assertTraces( + trace( + SORT_BY_START_TIME, + span().root().operationName("parent"), + span() + .childOfPrevious() + .operationName(H2_QUERY) + .resourceName("INSERT INTO peer_test (id, name) VALUES (1, 'test')") + .type(DDSpanTypes.SQL) + .measured() + .tags( + tag(Tags.COMPONENT, eqs("r2dbc")), + tag(Tags.SPAN_KIND, eqs(Tags.SPAN_KIND_CLIENT)), + tag(Tags.DB_TYPE, eqs("h2")), + tag(Tags.DB_INSTANCE, eqs("peerdb")), + tag(Tags.PEER_HOSTNAME, eqs("db.example.com")), + tag(Tags.DB_USER, eqs("testuser")), + tag("_dd.svc_src", any()), + defaultTags()))); + } + + @Test + void peerHostnameSetAcrossMultipleQueries() { + AgentSpan parent = startSpan("test", "parent"); + try (AgentScope scope = activateSpan(parent)) { + // First query: INSERT + Mono.from( + connection + .createStatement("INSERT INTO peer_test (id, name) VALUES (1, 'first')") + .execute()) + .flatMapMany(Result::getRowsUpdated) + .blockLast(); + + // Second query: SELECT + Flux.from(connection.createStatement("SELECT * FROM peer_test").execute()) + .flatMap(result -> result.map((row, metadata) -> row.get(0))) + .collectList() + .block(); + } finally { + parent.finish(); + } + + assertTraces( + trace( + SORT_BY_START_TIME, + span().root().operationName("parent"), + span() + .childOfPrevious() + .operationName(H2_QUERY) + .resourceName("INSERT INTO peer_test (id, name) VALUES (1, 'first')") + .type(DDSpanTypes.SQL) + .measured() + .tags( + tag(Tags.COMPONENT, eqs("r2dbc")), + tag(Tags.SPAN_KIND, eqs(Tags.SPAN_KIND_CLIENT)), + tag(Tags.DB_TYPE, eqs("h2")), + tag(Tags.DB_INSTANCE, eqs("peerdb")), + tag(Tags.PEER_HOSTNAME, eqs("db.example.com")), + tag(Tags.DB_USER, eqs("testuser")), + tag("_dd.svc_src", any()), + defaultTags()), + span() + .childOfIndex(0) + .operationName(H2_QUERY) + .resourceName("SELECT * FROM peer_test") + .type(DDSpanTypes.SQL) + .measured() + .tags( + tag(Tags.COMPONENT, eqs("r2dbc")), + tag(Tags.SPAN_KIND, eqs(Tags.SPAN_KIND_CLIENT)), + tag(Tags.DB_TYPE, eqs("h2")), + tag(Tags.DB_INSTANCE, eqs("peerdb")), + tag(Tags.PEER_HOSTNAME, eqs("db.example.com")), + tag(Tags.DB_USER, eqs("testuser")), + tag("_dd.svc_src", any()), + defaultTags()))); + } + + @Test + void peerHostnamePreservedOnErrorQuery() { + AgentSpan parent = startSpan("test", "parent"); + try (AgentScope scope = activateSpan(parent)) { + try { + Flux.from(connection.createStatement("SELECT * FROM nonexistent_peer_table").execute()) + .flatMap(result -> result.map((row, metadata) -> row.get(0))) + .collectList() + .block(); + } catch (Exception ignored) { + // Expected to fail + } + } finally { + parent.finish(); + } + + // Even on error, peer.hostname, db.instance, and db.user should still be set + assertTraces( + trace( + SORT_BY_START_TIME, + span().root().operationName("parent"), + span() + .childOfPrevious() + .operationName(H2_QUERY) + .resourceName("SELECT * FROM nonexistent_peer_table") + .type(DDSpanTypes.SQL) + .error() + .measured() + .tags( + tag(Tags.COMPONENT, eqs("r2dbc")), + tag(Tags.SPAN_KIND, eqs(Tags.SPAN_KIND_CLIENT)), + tag(Tags.DB_TYPE, eqs("h2")), + tag(Tags.DB_INSTANCE, eqs("peerdb")), + tag(Tags.PEER_HOSTNAME, eqs("db.example.com")), + tag(Tags.DB_USER, eqs("testuser")), + tag("_dd.svc_src", any()), + tag(DDTags.ERROR_MSG, any()), + error(Exception.class), + defaultTags()))); + } +} diff --git a/metadata/supported-configurations.json b/metadata/supported-configurations.json index 6c3fe354b68..7ea21a952b0 100644 --- a/metadata/supported-configurations.json +++ b/metadata/supported-configurations.json @@ -9324,6 +9324,30 @@ "aliases": ["DD_TRACE_INTEGRATION_RATPACK_REQUEST_BODY_ENABLED", "DD_INTEGRATION_RATPACK_REQUEST_BODY_ENABLED"] } ], + "DD_TRACE_R2DBC_ANALYTICS_ENABLED": [ + { + "version": "A", + "type": "boolean", + "default": "false", + "aliases": ["DD_R2DBC_ANALYTICS_ENABLED"] + } + ], + "DD_TRACE_R2DBC_ANALYTICS_SAMPLE_RATE": [ + { + "version": "A", + "type": "decimal", + "default": "1.0", + "aliases": ["DD_R2DBC_ANALYTICS_SAMPLE_RATE"] + } + ], + "DD_TRACE_R2DBC_ENABLED": [ + { + "version": "A", + "type": "boolean", + "default": "true", + "aliases": ["DD_TRACE_INTEGRATION_R2DBC_ENABLED", "DD_INTEGRATION_R2DBC_ENABLED"] + } + ], "DD_TRACE_REACTIVE_STREAMS_1_ENABLED": [ { "version": "A", diff --git a/settings.gradle.kts b/settings.gradle.kts index 14e988af3fe..a54ecf1edfd 100644 --- a/settings.gradle.kts +++ b/settings.gradle.kts @@ -538,6 +538,7 @@ include( ":dd-java-agent:instrumentation:quartz-2.0", ":dd-java-agent:instrumentation:rabbitmq-amqp-2.7", ":dd-java-agent:instrumentation:ratpack-1.5", + ":dd-java-agent:instrumentation:r2dbc:r2dbc-1.0", ":dd-java-agent:instrumentation:reactive-streams-1.0", ":dd-java-agent:instrumentation:reactor-core-3.1", ":dd-java-agent:instrumentation:reactor-netty-1.0",