diff --git a/dd-java-agent/agent-installer/build.gradle b/dd-java-agent/agent-installer/build.gradle index 00bfa0e4af5..ec176441bf6 100644 --- a/dd-java-agent/agent-installer/build.gradle +++ b/dd-java-agent/agent-installer/build.gradle @@ -31,6 +31,7 @@ dependencies { compileOnly project(':products:metrics:metrics-lib') testImplementation project(':dd-java-agent:testing') + testImplementation project(':utils:test-junit-utils') } tasks.named("compileMain_java11Java", JavaCompile) { diff --git a/dd-java-agent/agent-installer/src/main/java/datadog/trace/agent/tooling/AgentInstaller.java b/dd-java-agent/agent-installer/src/main/java/datadog/trace/agent/tooling/AgentInstaller.java index 3a8c7065362..d4fd57d270d 100644 --- a/dd-java-agent/agent-installer/src/main/java/datadog/trace/agent/tooling/AgentInstaller.java +++ b/dd-java-agent/agent-installer/src/main/java/datadog/trace/agent/tooling/AgentInstaller.java @@ -331,6 +331,9 @@ public static Set getEnabledSystems() { if (cfg.isUsmEnabled()) { enabledSystems.add(InstrumenterModule.TargetSystem.USM); } + if (cfg.isDataStreamsEnabled()) { + enabledSystems.add(InstrumenterModule.TargetSystem.DATA_STREAMS); + } if (cfg.isLlmObsEnabled()) { enabledSystems.add(InstrumenterModule.TargetSystem.LLMOBS); } diff --git a/dd-java-agent/agent-installer/src/test/java/datadog/trace/agent/tooling/AgentInstallerGetEnabledSystemsTest.java b/dd-java-agent/agent-installer/src/test/java/datadog/trace/agent/tooling/AgentInstallerGetEnabledSystemsTest.java new file mode 100644 index 00000000000..8352fc8d3cb --- /dev/null +++ b/dd-java-agent/agent-installer/src/test/java/datadog/trace/agent/tooling/AgentInstallerGetEnabledSystemsTest.java @@ -0,0 +1,75 @@ +package datadog.trace.agent.tooling; + +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import datadog.trace.agent.tooling.InstrumenterModule.TargetSystem; +import datadog.trace.test.junit.utils.config.WithConfig; +import datadog.trace.test.junit.utils.config.WithConfigExtension; +import java.util.Set; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; + +/** + * Tests for {@link AgentInstaller#getEnabledSystems()} to verify that it correctly includes target + * systems based on their corresponding configuration flags. + */ +@ExtendWith(WithConfigExtension.class) +class AgentInstallerGetEnabledSystemsTest { + + /** + * Verifies that DATA_STREAMS target system is not included when data.streams.enabled is false + * (default). + */ + @Test + void dataStreamsNotIncludedWhenDisabled() { + Set enabledSystems = AgentInstaller.getEnabledSystems(); + assertFalse( + enabledSystems.contains(TargetSystem.DATA_STREAMS), + "DATA_STREAMS should not be included when disabled"); + } + + /** Verifies that DATA_STREAMS target system is included when data.streams.enabled is true. */ + @Test + @WithConfig(key = "data.streams.enabled", value = "true") + void dataStreamsIncludedWhenEnabled() { + Set enabledSystems = AgentInstaller.getEnabledSystems(); + assertTrue( + enabledSystems.contains(TargetSystem.DATA_STREAMS), + "DATA_STREAMS should be included when enabled"); + } + + /** Verifies that USM target system is not included when usm.enabled is false (default). */ + @Test + void usmNotIncludedWhenDisabled() { + Set enabledSystems = AgentInstaller.getEnabledSystems(); + assertFalse( + enabledSystems.contains(TargetSystem.USM), "USM should not be included when disabled"); + } + + /** Verifies that USM target system is included when usm.enabled is true. */ + @Test + @WithConfig(key = "usm.enabled", value = "true") + void usmIncludedWhenEnabled() { + Set enabledSystems = AgentInstaller.getEnabledSystems(); + assertTrue(enabledSystems.contains(TargetSystem.USM), "USM should be included when enabled"); + } + + /** Verifies that LLMOBS target system is not included when llmobs.enabled is false (default). */ + @Test + void llmobsNotIncludedWhenDisabled() { + Set enabledSystems = AgentInstaller.getEnabledSystems(); + assertFalse( + enabledSystems.contains(TargetSystem.LLMOBS), + "LLMOBS should not be included when disabled"); + } + + /** Verifies that LLMOBS target system is included when llmobs.enabled is true. */ + @Test + @WithConfig(key = "llmobs.enabled", value = "true") + void llmobsIncludedWhenEnabled() { + Set enabledSystems = AgentInstaller.getEnabledSystems(); + assertTrue( + enabledSystems.contains(TargetSystem.LLMOBS), "LLMOBS should be included when enabled"); + } +} diff --git a/dd-java-agent/agent-tooling/build.gradle b/dd-java-agent/agent-tooling/build.gradle index 06272944f25..c7a0694523b 100644 --- a/dd-java-agent/agent-tooling/build.gradle +++ b/dd-java-agent/agent-tooling/build.gradle @@ -49,6 +49,7 @@ dependencies { api libs.bytebuddyagent testImplementation project(':dd-java-agent:testing') + testImplementation project(':utils:test-junit-utils') testImplementation libs.bytebuddy testImplementation group: 'com.google.guava', name: 'guava-testlib', version: '20.0' diff --git a/dd-java-agent/agent-tooling/src/main/java/datadog/trace/agent/tooling/InstrumenterModule.java b/dd-java-agent/agent-tooling/src/main/java/datadog/trace/agent/tooling/InstrumenterModule.java index d2abbc265e5..57e11f6bcfe 100644 --- a/dd-java-agent/agent-tooling/src/main/java/datadog/trace/agent/tooling/InstrumenterModule.java +++ b/dd-java-agent/agent-tooling/src/main/java/datadog/trace/agent/tooling/InstrumenterModule.java @@ -41,6 +41,7 @@ public abstract class InstrumenterModule implements Instrumenter { *
  • {@link TargetSystem#IAST iast} *
  • {@link TargetSystem#CIVISIBILITY ci-visibility} *
  • {@link TargetSystem#USM usm} + *
  • {@link TargetSystem#DATA_STREAMS data-streams} *
  • {@link TargetSystem#CONTEXT_TRACKING context-tracking} *
  • {@link TargetSystem#RASP rasp} * @@ -53,6 +54,7 @@ public enum TargetSystem { CIVISIBILITY, USM, LLMOBS, + DATA_STREAMS, CONTEXT_TRACKING, RASP, } @@ -320,6 +322,24 @@ public final boolean isApplicable(Set enabledSystems) { } } + /** Parent class for instrumentations that support both tracing and Data Streams Monitoring */ + public abstract static class DataStreams extends InstrumenterModule { + public DataStreams(String instrumentationName, String... additionalNames) { + super(instrumentationName, additionalNames); + } + + @Override + public final boolean isApplicable(Set enabledSystems) { + return enabledSystems.contains(TargetSystem.TRACING) + || enabledSystems.contains(TargetSystem.DATA_STREAMS); + } + + @Override + public boolean isEnabled() { + return super.isEnabled() || InstrumenterConfig.get().isDataStreamsEnabled(); + } + } + /** Parent class for all CI related instrumentations */ public abstract static class CiVisibility extends InstrumenterModule { public CiVisibility(String instrumentationName, String... additionalNames) { diff --git a/dd-java-agent/agent-tooling/src/test/java/datadog/trace/agent/tooling/InstrumenterModuleTest.java b/dd-java-agent/agent-tooling/src/test/java/datadog/trace/agent/tooling/InstrumenterModuleTest.java new file mode 100644 index 00000000000..1fb55ea9103 --- /dev/null +++ b/dd-java-agent/agent-tooling/src/test/java/datadog/trace/agent/tooling/InstrumenterModuleTest.java @@ -0,0 +1,102 @@ +package datadog.trace.agent.tooling; + +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import datadog.trace.agent.tooling.InstrumenterModule.TargetSystem; +import datadog.trace.test.junit.utils.config.WithConfig; +import datadog.trace.test.junit.utils.config.WithConfigExtension; +import java.util.HashSet; +import java.util.Set; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; + +@ExtendWith(WithConfigExtension.class) +class InstrumenterModuleTest { + + @Test + void testDataStreamsIsApplicableWithTracing() { + Set enabledSystems = new HashSet<>(); + enabledSystems.add(TargetSystem.TRACING); + + InstrumenterModule.DataStreams module = new InstrumenterModule.DataStreams("test-module") {}; + + assertTrue(module.isApplicable(enabledSystems)); + } + + @Test + void testDataStreamsIsApplicableWithDataStreams() { + Set enabledSystems = new HashSet<>(); + enabledSystems.add(TargetSystem.DATA_STREAMS); + + InstrumenterModule.DataStreams module = new InstrumenterModule.DataStreams("test-module") {}; + + assertTrue(module.isApplicable(enabledSystems)); + } + + @Test + void testDataStreamsIsApplicableWithBoth() { + Set enabledSystems = new HashSet<>(); + enabledSystems.add(TargetSystem.TRACING); + enabledSystems.add(TargetSystem.DATA_STREAMS); + + InstrumenterModule.DataStreams module = new InstrumenterModule.DataStreams("test-module") {}; + + assertTrue(module.isApplicable(enabledSystems)); + } + + @Test + void testDataStreamsIsApplicableWithNeither() { + Set enabledSystems = new HashSet<>(); + enabledSystems.add(TargetSystem.APPSEC); + + InstrumenterModule.DataStreams module = new InstrumenterModule.DataStreams("test-module") {}; + + assertFalse(module.isApplicable(enabledSystems)); + } + + @Test + @WithConfig(key = "trace.test-kafka-module.enabled", value = "false") + @WithConfig(key = "data.streams.enabled", value = "true") + void testDataStreamsIsEnabledWhenDataStreamsEnabledOverridesFalse() { + // When tracing for this integration is disabled but DSM is explicitly enabled, + // isEnabled() should still return true. + InstrumenterModule.DataStreams module = + new InstrumenterModule.DataStreams("test-kafka-module") {}; + + assertTrue(module.isEnabled()); + } + + @Test + @WithConfig(key = "trace.test-kafka-module.enabled", value = "true") + @WithConfig(key = "data.streams.enabled", value = "false") + void testDataStreamsIsEnabledWhenSuperEnabledIsTrue() { + // When super.isEnabled() is true, isEnabled() should return true regardless of DSM state. + InstrumenterModule.DataStreams module = + new InstrumenterModule.DataStreams("test-kafka-module") {}; + + assertTrue(module.isEnabled()); + } + + @Test + @WithConfig(key = "trace.test-kafka-module.enabled", value = "true") + @WithConfig(key = "data.streams.enabled", value = "true") + void testDataStreamsIsEnabledWhenBothEnabled() { + // When both super.isEnabled() and DSM are enabled, isEnabled() should return true. + InstrumenterModule.DataStreams module = + new InstrumenterModule.DataStreams("test-kafka-module") {}; + + assertTrue(module.isEnabled()); + } + + @Test + @WithConfig(key = "trace.test-kafka-module.enabled", value = "false") + @WithConfig(key = "data.streams.enabled", value = "false") + void testDataStreamsIsEnabledWhenBothDisabled() { + // When both super.isEnabled() and DSM are disabled, isEnabled() should return false. + InstrumenterModule.DataStreams module = + new InstrumenterModule.DataStreams("test-kafka-module") {}; + + assertFalse(module.isEnabled()); + } +} diff --git a/dd-java-agent/instrumentation/kafka/kafka-clients-0.11/src/main/java/datadog/trace/instrumentation/kafka_clients/ConsumerCoordinatorInstrumentation.java b/dd-java-agent/instrumentation/kafka/kafka-clients-0.11/src/main/java/datadog/trace/instrumentation/kafka_clients/ConsumerCoordinatorInstrumentation.java index 57de8c42ff3..962730d5a22 100644 --- a/dd-java-agent/instrumentation/kafka/kafka-clients-0.11/src/main/java/datadog/trace/instrumentation/kafka_clients/ConsumerCoordinatorInstrumentation.java +++ b/dd-java-agent/instrumentation/kafka/kafka-clients-0.11/src/main/java/datadog/trace/instrumentation/kafka_clients/ConsumerCoordinatorInstrumentation.java @@ -26,11 +26,11 @@ import org.apache.kafka.common.TopicPartition; @AutoService(InstrumenterModule.class) -public final class ConsumerCoordinatorInstrumentation extends InstrumenterModule.Tracing +public final class ConsumerCoordinatorInstrumentation extends InstrumenterModule.DataStreams implements Instrumenter.ForSingleType, Instrumenter.HasMethodAdvice { public ConsumerCoordinatorInstrumentation() { - super("kafka", "kafka-0.11"); + super(KafkaDecorator.INTEGRATION_NAME, KafkaDecorator.LEGACY_INTEGRATION_NAME); } @Override diff --git a/dd-java-agent/instrumentation/kafka/kafka-clients-0.11/src/main/java/datadog/trace/instrumentation/kafka_clients/KafkaConsumerInfoInstrumentation.java b/dd-java-agent/instrumentation/kafka/kafka-clients-0.11/src/main/java/datadog/trace/instrumentation/kafka_clients/KafkaConsumerInfoInstrumentation.java index a2fa481491c..9446334ce24 100644 --- a/dd-java-agent/instrumentation/kafka/kafka-clients-0.11/src/main/java/datadog/trace/instrumentation/kafka_clients/KafkaConsumerInfoInstrumentation.java +++ b/dd-java-agent/instrumentation/kafka/kafka-clients-0.11/src/main/java/datadog/trace/instrumentation/kafka_clients/KafkaConsumerInfoInstrumentation.java @@ -44,11 +44,11 @@ * and cluster ID, in the context store for later use. */ @AutoService(InstrumenterModule.class) -public final class KafkaConsumerInfoInstrumentation extends InstrumenterModule.Tracing +public final class KafkaConsumerInfoInstrumentation extends InstrumenterModule.DataStreams implements Instrumenter.ForSingleType, Instrumenter.HasMethodAdvice { public KafkaConsumerInfoInstrumentation() { - super("kafka", "kafka-0.11"); + super(KafkaDecorator.INTEGRATION_NAME, KafkaDecorator.LEGACY_INTEGRATION_NAME); } @Override @@ -252,7 +252,10 @@ public static AgentScope onEnter(@Advice.This KafkaConsumer consumer) { } } - if (traceConfig().isDataStreamsEnabled()) { + if (traceConfig().isDataStreamsEnabled() && KafkaDecorator.TRACING_ENABLED) { + // DSM-only mode (tracing disabled) never creates a real poll span: TracingIterator + // carries its pathway context on a lightweight, never-collected span shim instead, so + // there's nothing here that needs wrapping/protecting from being force-dropped. final AgentSpan span = startSpan(JAVA_KAFKA.toString(), KAFKA_POLL); return activateSpan(span); } diff --git a/dd-java-agent/instrumentation/kafka/kafka-clients-0.11/src/main/java/datadog/trace/instrumentation/kafka_clients/KafkaConsumerInstrumentation.java b/dd-java-agent/instrumentation/kafka/kafka-clients-0.11/src/main/java/datadog/trace/instrumentation/kafka_clients/KafkaConsumerInstrumentation.java index 756f59aad4b..d516d48a169 100644 --- a/dd-java-agent/instrumentation/kafka/kafka-clients-0.11/src/main/java/datadog/trace/instrumentation/kafka_clients/KafkaConsumerInstrumentation.java +++ b/dd-java-agent/instrumentation/kafka/kafka-clients-0.11/src/main/java/datadog/trace/instrumentation/kafka_clients/KafkaConsumerInstrumentation.java @@ -27,11 +27,11 @@ import org.apache.kafka.clients.consumer.ConsumerRecords; @AutoService(InstrumenterModule.class) -public final class KafkaConsumerInstrumentation extends InstrumenterModule.Tracing +public final class KafkaConsumerInstrumentation extends InstrumenterModule.DataStreams implements Instrumenter.ForSingleType, Instrumenter.HasMethodAdvice { public KafkaConsumerInstrumentation() { - super("kafka", "kafka-0.11"); + super(KafkaDecorator.INTEGRATION_NAME, KafkaDecorator.LEGACY_INTEGRATION_NAME); } @Override diff --git a/dd-java-agent/instrumentation/kafka/kafka-clients-0.11/src/main/java/datadog/trace/instrumentation/kafka_clients/KafkaDecorator.java b/dd-java-agent/instrumentation/kafka/kafka-clients-0.11/src/main/java/datadog/trace/instrumentation/kafka_clients/KafkaDecorator.java index 53c579a4dee..eef3c71103e 100644 --- a/dd-java-agent/instrumentation/kafka/kafka-clients-0.11/src/main/java/datadog/trace/instrumentation/kafka_clients/KafkaDecorator.java +++ b/dd-java-agent/instrumentation/kafka/kafka-clients-0.11/src/main/java/datadog/trace/instrumentation/kafka_clients/KafkaDecorator.java @@ -12,6 +12,7 @@ import datadog.trace.api.Config; import datadog.trace.api.Functions; +import datadog.trace.api.InstrumenterConfig; import datadog.trace.api.cache.DDCache; import datadog.trace.api.cache.DDCaches; import datadog.trace.api.naming.SpanNaming; @@ -20,6 +21,7 @@ import datadog.trace.bootstrap.instrumentation.api.Tags; import datadog.trace.bootstrap.instrumentation.api.UTF8BytesString; import datadog.trace.bootstrap.instrumentation.decorator.MessagingClientDecorator; +import java.util.Arrays; import java.util.function.Function; import java.util.function.Supplier; import org.apache.kafka.clients.consumer.ConsumerRecord; @@ -29,6 +31,10 @@ public class KafkaDecorator extends MessagingClientDecorator { private static final String KAFKA = "kafka"; + // Kept in sync with the names each kafka-clients-0.11 instrumentation module passes to its own + // super(...) constructor call, so TRACING_ENABLED can't drift from what is actually registered. + public static final String INTEGRATION_NAME = KAFKA; + public static final String LEGACY_INTEGRATION_NAME = "kafka-0.11"; public static final CharSequence JAVA_KAFKA = UTF8BytesString.create("java-kafka"); public static final CharSequence KAFKA_CONSUME = UTF8BytesString.create( @@ -42,6 +48,11 @@ public class KafkaDecorator extends MessagingClientDecorator { public static final boolean KAFKA_LEGACY_TRACING = Config.get().isKafkaLegacyTracingEnabled(); public static final boolean TIME_IN_QUEUE_ENABLED = Config.get().isTimeInQueueEnabled(!KAFKA_LEGACY_TRACING, KAFKA); + public static final boolean TRACING_ENABLED = + InstrumenterConfig.get() + .isIntegrationEnabled( + Arrays.asList(INTEGRATION_NAME, LEGACY_INTEGRATION_NAME), + InstrumenterConfig.get().isIntegrationsEnabled()); public static final String KAFKA_PRODUCED_KEY = "x_datadog_kafka_produced"; private final String spanKind; private final CharSequence spanType; diff --git a/dd-java-agent/instrumentation/kafka/kafka-clients-0.11/src/main/java/datadog/trace/instrumentation/kafka_clients/KafkaProducerInstrumentation.java b/dd-java-agent/instrumentation/kafka/kafka-clients-0.11/src/main/java/datadog/trace/instrumentation/kafka_clients/KafkaProducerInstrumentation.java index 227d8872648..580ce12d19b 100644 --- a/dd-java-agent/instrumentation/kafka/kafka-clients-0.11/src/main/java/datadog/trace/instrumentation/kafka_clients/KafkaProducerInstrumentation.java +++ b/dd-java-agent/instrumentation/kafka/kafka-clients-0.11/src/main/java/datadog/trace/instrumentation/kafka_clients/KafkaProducerInstrumentation.java @@ -12,12 +12,14 @@ import static datadog.trace.bootstrap.instrumentation.api.AgentTracer.activateSpan; import static datadog.trace.bootstrap.instrumentation.api.AgentTracer.activeSpan; import static datadog.trace.bootstrap.instrumentation.api.AgentTracer.startSpan; +import static datadog.trace.bootstrap.instrumentation.api.AgentTracer.traceConfig; import static datadog.trace.instrumentation.kafka_clients.KafkaDecorator.JAVA_KAFKA; import static datadog.trace.instrumentation.kafka_clients.KafkaDecorator.KAFKA_PRODUCE; import static datadog.trace.instrumentation.kafka_clients.KafkaDecorator.PRODUCER_DECORATE; import static datadog.trace.instrumentation.kafka_clients.KafkaDecorator.TIME_IN_QUEUE_ENABLED; import static datadog.trace.instrumentation.kafka_common.StreamingContext.STREAMING_CONTEXT; import static datadog.trace.instrumentation.kafka_common.Utils.DSM_TRANSACTION_SOURCE_READER; +import static datadog.trace.instrumentation.kafka_common.Utils.newPathwayOnlySpan; import static java.util.Collections.singletonMap; import static net.bytebuddy.matcher.ElementMatchers.isConstructor; import static net.bytebuddy.matcher.ElementMatchers.isMethod; @@ -58,11 +60,11 @@ import org.apache.kafka.common.record.RecordBatch; @AutoService(InstrumenterModule.class) -public final class KafkaProducerInstrumentation extends InstrumenterModule.Tracing +public final class KafkaProducerInstrumentation extends InstrumenterModule.DataStreams implements Instrumenter.ForSingleType, Instrumenter.HasMethodAdvice { public KafkaProducerInstrumentation() { - super("kafka", "kafka-0.11"); + super(KafkaDecorator.INTEGRATION_NAME, KafkaDecorator.LEGACY_INTEGRATION_NAME); } @Override @@ -154,20 +156,32 @@ public static AgentScope onEnter( final AgentSpanContext extractedContext = extractContextAndGetSpanContext(record.headers(), TextMapExtractAdapter.GETTER); - final AgentSpan localActiveSpan = activeSpan(); - final AgentSpan span; final AgentSpan callbackParentSpan; - if (extractedContext != null) { - span = startSpan(JAVA_KAFKA.toString(), KAFKA_PRODUCE, extractedContext); + if (!KafkaDecorator.TRACING_ENABLED && traceConfig().isDataStreamsEnabled()) { + // DSM-only mode: never create a real span, so nothing for this integration is ever + // written to the agent. The pathway is carried on a lightweight, never-collected span + // shim instead, falling back to whatever pathway the currently active span (if any) + // is carrying so a consume->produce chain keeps propagating the same pathway. + final AgentSpan localActiveSpan = activeSpan(); + final AgentSpanContext pathwaySource = + extractedContext != null + ? extractedContext + : localActiveSpan == null ? null : localActiveSpan.spanContext(); + span = newPathwayOnlySpan(pathwaySource); callbackParentSpan = span; } else { - span = startSpan(JAVA_KAFKA.toString(), KAFKA_PRODUCE); - callbackParentSpan = localActiveSpan; + if (extractedContext != null) { + span = startSpan(JAVA_KAFKA.toString(), KAFKA_PRODUCE, extractedContext); + callbackParentSpan = span; + } else { + span = startSpan(JAVA_KAFKA.toString(), KAFKA_PRODUCE); + callbackParentSpan = activeSpan(); + } + PRODUCER_DECORATE.afterStart(span); + PRODUCER_DECORATE.onProduce(span, record, producerConfig, clusterId); } - PRODUCER_DECORATE.afterStart(span); - PRODUCER_DECORATE.onProduce(span, record, producerConfig, clusterId); callback = new KafkaProducerCallback(callback, callbackParentSpan, span, clusterId); diff --git a/dd-java-agent/instrumentation/kafka/kafka-clients-0.11/src/main/java/datadog/trace/instrumentation/kafka_clients/MetadataInstrumentation.java b/dd-java-agent/instrumentation/kafka/kafka-clients-0.11/src/main/java/datadog/trace/instrumentation/kafka_clients/MetadataInstrumentation.java index d6acfe30369..871e82cc6b5 100644 --- a/dd-java-agent/instrumentation/kafka/kafka-clients-0.11/src/main/java/datadog/trace/instrumentation/kafka_clients/MetadataInstrumentation.java +++ b/dd-java-agent/instrumentation/kafka/kafka-clients-0.11/src/main/java/datadog/trace/instrumentation/kafka_clients/MetadataInstrumentation.java @@ -24,11 +24,11 @@ import org.apache.kafka.common.requests.MetadataResponse; @AutoService(InstrumenterModule.class) -public class MetadataInstrumentation extends InstrumenterModule.Tracing +public class MetadataInstrumentation extends InstrumenterModule.DataStreams implements Instrumenter.ForTypeHierarchy, Instrumenter.HasMethodAdvice { public MetadataInstrumentation() { - super("kafka", "kafka-0.11"); + super(KafkaDecorator.INTEGRATION_NAME, KafkaDecorator.LEGACY_INTEGRATION_NAME); } @Override diff --git a/dd-java-agent/instrumentation/kafka/kafka-clients-0.11/src/main/java/datadog/trace/instrumentation/kafka_clients/TracingIterator.java b/dd-java-agent/instrumentation/kafka/kafka-clients-0.11/src/main/java/datadog/trace/instrumentation/kafka_clients/TracingIterator.java index 93179e4e3f2..c53d56a4ea5 100644 --- a/dd-java-agent/instrumentation/kafka/kafka-clients-0.11/src/main/java/datadog/trace/instrumentation/kafka_clients/TracingIterator.java +++ b/dd-java-agent/instrumentation/kafka/kafka-clients-0.11/src/main/java/datadog/trace/instrumentation/kafka_clients/TracingIterator.java @@ -18,6 +18,7 @@ import static datadog.trace.instrumentation.kafka_common.StreamingContext.STREAMING_CONTEXT; import static datadog.trace.instrumentation.kafka_common.Utils.DSM_TRANSACTION_SOURCE_READER; import static datadog.trace.instrumentation.kafka_common.Utils.computePayloadSizeBytes; +import static datadog.trace.instrumentation.kafka_common.Utils.newPathwayOnlySpan; import static java.util.concurrent.TimeUnit.MILLISECONDS; import datadog.context.Context; @@ -97,56 +98,11 @@ protected void startNewRecordSpan(ConsumerRecord val) { previousSpan.finishWithEndToEnd(); } } - AgentSpan span, queueSpan = null; if (val != null) { - if (!Config.get().isKafkaClientPropagationDisabledForTopic(val.topic())) { - final AgentSpanContext spanContext = - extractContextAndGetSpanContext(val.headers(), GETTER); - long timeInQueueStart = GETTER.extractTimeInQueueStart(val.headers()); - if (timeInQueueStart == 0 || !TIME_IN_QUEUE_ENABLED) { - span = startSpan(JAVA_KAFKA.toString(), operationName, spanContext); - } else { - queueSpan = - startSpan( - JAVA_KAFKA.toString(), - KAFKA_DELIVER, - spanContext, - MILLISECONDS.toMicros(timeInQueueStart)); - BROKER_DECORATE.afterStart(queueSpan); - BROKER_DECORATE.onTimeInQueue(queueSpan, val); - span = startSpan(JAVA_KAFKA.toString(), operationName, queueSpan.spanContext()); - BROKER_DECORATE.beforeFinish(queueSpan); - // The queueSpan will be finished after inner span has been activated to ensure that - // spans are written out together by TraceStructureWriter when running in strict mode - } - - DataStreamsTags tags = create("kafka", INBOUND, val.topic(), group, clusterId); - final long payloadSize = - traceConfig().isDataStreamsEnabled() ? computePayloadSizeBytes(val) : 0; - if (STREAMING_CONTEXT.isDisabledForTopic(val.topic())) { - AgentTracer.get() - .getDataStreamsMonitoring() - .setCheckpoint(span, create(tags, val.timestamp(), payloadSize)); - } else { - // when we're in a streaming context we want to consume only from source topics - if (STREAMING_CONTEXT.isSourceTopic(val.topic())) { - // We have to inject the context to headers here, - // since the data received from the source may leave the topology on - // some other instance of the application, breaking the context propagation - // for DSM users - Propagator dsmPropagator = Propagators.forConcern(DSM_CONCERN); - DataStreamsContext dsmContext = create(tags, val.timestamp(), payloadSize); - dsmPropagator.inject(span.with(dsmContext), val.headers(), SETTER); - } - } - } else { - span = startSpan(JAVA_KAFKA.toString(), operationName, null); - } - if (val.value() == null) { - span.setTag(InstrumentationTags.TOMBSTONE, true); - } - decorator.afterStart(span); - decorator.onConsume(span, val, group, clusterId, bootstrapServers); + final AgentSpan span = + !KafkaDecorator.TRACING_ENABLED && traceConfig().isDataStreamsEnabled() + ? startDsmOnlyPathwaySpan(val) + : startTracedConsumeSpan(val); if (InstrumenterConfig.get().isLegacyContextManagerEnabled()) { activateNext(span); } else { @@ -155,23 +111,108 @@ protected void startNewRecordSpan(ConsumerRecord val) { previousSpan.finishWithEndToEnd(); } } - if (null != queueSpan) { - queueSpan.finish(); - } - - AgentTracer.get() - .getDataStreamsMonitoring() - .trackTransaction( - span, - DataStreamsTransactionExtractor.Type.KAFKA_CONSUME_HEADERS, - val.headers(), - DSM_TRANSACTION_SOURCE_READER); } } catch (final Exception e) { log.debug("Error starting new record span", e); } } + /** + * Creates and activates the real APM consume span (and, when time-in-queue is enabled, its broker + * parent), tags it, and reports DSM checkpoints/transactions off of it. + */ + private AgentSpan startTracedConsumeSpan(ConsumerRecord val) { + AgentSpan span, queueSpan = null; + if (!Config.get().isKafkaClientPropagationDisabledForTopic(val.topic())) { + final AgentSpanContext spanContext = extractContextAndGetSpanContext(val.headers(), GETTER); + long timeInQueueStart = GETTER.extractTimeInQueueStart(val.headers()); + if (timeInQueueStart == 0 || !TIME_IN_QUEUE_ENABLED) { + span = startSpan(JAVA_KAFKA.toString(), operationName, spanContext); + } else { + queueSpan = + startSpan( + JAVA_KAFKA.toString(), + KAFKA_DELIVER, + spanContext, + MILLISECONDS.toMicros(timeInQueueStart)); + BROKER_DECORATE.afterStart(queueSpan); + BROKER_DECORATE.onTimeInQueue(queueSpan, val); + span = startSpan(JAVA_KAFKA.toString(), operationName, queueSpan.spanContext()); + BROKER_DECORATE.beforeFinish(queueSpan); + // The queueSpan will be finished after inner span has been activated to ensure that + // spans are written out together by TraceStructureWriter when running in strict mode + } + + DataStreamsTags tags = create("kafka", INBOUND, val.topic(), group, clusterId); + final long payloadSize = + traceConfig().isDataStreamsEnabled() ? computePayloadSizeBytes(val) : 0; + reportDsmCheckpointOrInject(span, val, tags, payloadSize); + } else { + span = startSpan(JAVA_KAFKA.toString(), operationName, null); + } + if (val.value() == null) { + span.setTag(InstrumentationTags.TOMBSTONE, true); + } + decorator.afterStart(span); + decorator.onConsume(span, val, group, clusterId, bootstrapServers); + if (null != queueSpan) { + queueSpan.finish(); + } + + trackDsmConsumeTransaction(span, val); + return span; + } + + /** + * DSM-only mode (tracing disabled for kafka, DSM enabled): never creates a real span, so no span + * is ever written to the agent for this integration. Only the pathway checkpoint/injection and + * transaction tracking happen, carried by a lightweight, never-collected span shim. + */ + private AgentSpan startDsmOnlyPathwaySpan(ConsumerRecord val) { + AgentSpan span; + if (!Config.get().isKafkaClientPropagationDisabledForTopic(val.topic())) { + final AgentSpanContext extractedContext = + extractContextAndGetSpanContext(val.headers(), GETTER); + span = newPathwayOnlySpan(extractedContext); + DataStreamsTags tags = create("kafka", INBOUND, val.topic(), group, clusterId); + final long payloadSize = computePayloadSizeBytes(val); + reportDsmCheckpointOrInject(span, val, tags, payloadSize); + } else { + span = newPathwayOnlySpan(null); + } + trackDsmConsumeTransaction(span, val); + return span; + } + + /** + * Reports a DSM checkpoint for {@code val}'s topic, or - when in a streaming context and {@code + * val}'s topic is a source topic - injects the pathway context into its headers so it survives + * leaving the topology on another instance of the application. + */ + private void reportDsmCheckpointOrInject( + AgentSpan span, ConsumerRecord val, DataStreamsTags tags, long payloadSize) { + if (STREAMING_CONTEXT.isDisabledForTopic(val.topic())) { + AgentTracer.get() + .getDataStreamsMonitoring() + .setCheckpoint(span, create(tags, val.timestamp(), payloadSize)); + } else if (STREAMING_CONTEXT.isSourceTopic(val.topic())) { + // when we're in a streaming context we want to consume only from source topics + Propagator dsmPropagator = Propagators.forConcern(DSM_CONCERN); + DataStreamsContext dsmContext = create(tags, val.timestamp(), payloadSize); + dsmPropagator.inject(span.with(dsmContext), val.headers(), SETTER); + } + } + + private void trackDsmConsumeTransaction(AgentSpan span, ConsumerRecord val) { + AgentTracer.get() + .getDataStreamsMonitoring() + .trackTransaction( + span, + DataStreamsTransactionExtractor.Type.KAFKA_CONSUME_HEADERS, + val.headers(), + DSM_TRANSACTION_SOURCE_READER); + } + @Override public void remove() { delegateIterator.remove(); diff --git a/dd-java-agent/instrumentation/kafka/kafka-clients-0.11/src/test/groovy/KafkaClientDataStreamsOnlyForkedTest.groovy b/dd-java-agent/instrumentation/kafka/kafka-clients-0.11/src/test/groovy/KafkaClientDataStreamsOnlyForkedTest.groovy new file mode 100644 index 00000000000..77a6b18c51a --- /dev/null +++ b/dd-java-agent/instrumentation/kafka/kafka-clients-0.11/src/test/groovy/KafkaClientDataStreamsOnlyForkedTest.groovy @@ -0,0 +1,253 @@ +import datadog.trace.agent.test.InstrumentationSpecification +import datadog.trace.common.writer.ListWriter +import datadog.trace.core.DDSpan +import org.apache.kafka.clients.consumer.ConsumerConfig +import org.apache.kafka.clients.consumer.KafkaConsumer +import org.apache.kafka.clients.producer.KafkaProducer +import org.apache.kafka.clients.producer.ProducerRecord +import org.apache.kafka.common.TopicPartition +import org.apache.kafka.common.header.internals.RecordHeader +import org.apache.kafka.common.header.internals.RecordHeaders +import org.apache.kafka.common.serialization.StringSerializer +import org.springframework.kafka.test.rule.KafkaEmbedded +import org.springframework.kafka.test.utils.KafkaTestUtils + +import java.nio.charset.StandardCharsets + +import static datadog.trace.agent.test.utils.TraceUtils.runUnderTrace + +/** + * DSM-only coverage for kafka-clients-0.11: when Kafka APM tracing is disabled + * (integrations.enabled=false, or the per-integration trace.kafka.enabled=false override) but DSM + * is enabled, this integration must never create or write a real span/trace, while pathway + * checkpoints must still be tracked. These specs deliberately extend + * {@code InstrumentationSpecification} directly rather than {@code KafkaClientTestBase}: that base + * carries many concrete tests which assert on APM spans that are correctly absent once Kafka + * tracing is off, and Spock would run them (and fail) as inherited tests on every subclass. + */ +abstract class KafkaClientDataStreamsOnlyForkedTest extends InstrumentationSpecification { + static final SHARED_TOPIC = "shared.topic" + + KafkaEmbedded embeddedKafka + + def setup() { + embeddedKafka = new KafkaEmbedded(1, true, SHARED_TOPIC) + embeddedKafka.before() + } + + def cleanup() { + embeddedKafka?.after() + } + + @Override + boolean useStrictTraceWrites() { + return false + } + + @Override + protected boolean isDataStreamsEnabled() { + return true + } + + protected KafkaProducer newProducer() { + return new KafkaProducer( + KafkaTestUtils.senderProps(embeddedKafka.getBrokersAsString()), + new StringSerializer(), + new StringSerializer()) + } +} + +/** + * Regression guard for the contract that "integration disabled => zero spans of that type ever + * reach the agent": producing and consuming a message in DSM-only mode must not write any trace, + * even though DSM checkpoints for the same produce/consume are still tracked. + */ +class KafkaClientDataStreamsOnlyLocalRootForkedTest extends KafkaClientDataStreamsOnlyForkedTest { + @Override + void configurePreAgent() { + super.configurePreAgent() + injectSysConfig("integrations.enabled", "false") + } + + def "local-root produce and consume write no spans when kafka tracing is disabled and DSM is enabled"() { + setup: + def kafkaPartition = 0 + def consumerProperties = KafkaTestUtils.consumerProps("sender", "false", embeddedKafka) + consumerProperties.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest") + def consumer = new KafkaConsumer(consumerProperties) + def producer = newProducer() + consumer.assign(Arrays.asList(new TopicPartition(SHARED_TOPIC, kafkaPartition))) + + when: "a message is produced with no propagated trace headers, i.e. a genuine local root" + producer.send(new ProducerRecord(SHARED_TOPIC, kafkaPartition, null, "local-root-message")).get() + def recs = KafkaTestUtils.getRecords(consumer) + .records(new TopicPartition(SHARED_TOPIC, kafkaPartition)).iterator() + + then: "the message is delivered and DSM checkpoints are tracked, but no span is ever written" + recs.hasNext() + recs.next().value() == "local-root-message" + !recs.hasNext() + TEST_DATA_STREAMS_WRITER.waitForGroups(2) + TEST_WRITER.isEmpty() + + cleanup: + consumer?.close() + producer?.close() + } +} + +/** + * Regression guard: a message carrying a real, externally-propagated Datadog trace context must + * also write no span, under the same "kafka tracing disabled + DSM enabled" configuration - the + * DSM-only decision is based purely on the tracing/DSM config flags, not on the record's headers. + */ +class KafkaClientDataStreamsOnlyExtractedParentForkedTest extends KafkaClientDataStreamsOnlyForkedTest { + @Override + void configurePreAgent() { + super.configurePreAgent() + injectSysConfig("integrations.enabled", "false") + } + + def "produce and consume with an extracted parent trace context still write no spans"() { + setup: + def kafkaPartition = 0 + def consumerProperties = KafkaTestUtils.consumerProps("sender", "false", embeddedKafka) + consumerProperties.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest") + def consumer = new KafkaConsumer(consumerProperties) + def producer = newProducer() + consumer.assign(Arrays.asList(new TopicPartition(SHARED_TOPIC, kafkaPartition))) + + def existingTraceId = 1234567890123456L + def existingSpanId = 9876543210987654L + def headers = new RecordHeaders() + headers.add(new RecordHeader("x-datadog-trace-id", + String.valueOf(existingTraceId).getBytes(StandardCharsets.UTF_8))) + headers.add(new RecordHeader("x-datadog-parent-id", + String.valueOf(existingSpanId).getBytes(StandardCharsets.UTF_8))) + + when: "a message carrying a real, externally-propagated Datadog trace context is produced" + producer.send(new ProducerRecord(SHARED_TOPIC, kafkaPartition, null, "propagated-trace-message", headers)).get() + def recs = KafkaTestUtils.getRecords(consumer) + .records(new TopicPartition(SHARED_TOPIC, kafkaPartition)).iterator() + + then: "the message is delivered and DSM checkpoints are tracked, but no span is ever written" + recs.hasNext() + recs.next().value() == "propagated-trace-message" + !recs.hasNext() + TEST_DATA_STREAMS_WRITER.waitForGroups(2) + TEST_WRITER.isEmpty() + + cleanup: + consumer?.close() + producer?.close() + } +} + +/** + * Confirms the per-integration override (trace.kafka.enabled=false) suppresses span creation + * identically to the global integrations.enabled=false toggle used above. + */ +class KafkaClientDataStreamsOnlyIntegrationOverrideForkedTest extends KafkaClientDataStreamsOnlyForkedTest { + @Override + void configurePreAgent() { + super.configurePreAgent() + injectSysConfig("trace.kafka.enabled", "false") + } + + def "local-root produce and consume write no spans when trace.kafka.enabled=false and DSM is enabled"() { + setup: + def kafkaPartition = 0 + def consumerProperties = KafkaTestUtils.consumerProps("sender", "false", embeddedKafka) + consumerProperties.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest") + def consumer = new KafkaConsumer(consumerProperties) + def producer = newProducer() + consumer.assign(Arrays.asList(new TopicPartition(SHARED_TOPIC, kafkaPartition))) + + when: "a message is produced with no propagated trace headers, i.e. a genuine local root" + producer.send(new ProducerRecord(SHARED_TOPIC, kafkaPartition, null, "local-root-message")).get() + def recs = KafkaTestUtils.getRecords(consumer) + .records(new TopicPartition(SHARED_TOPIC, kafkaPartition)).iterator() + + then: "the message is delivered and DSM checkpoints are tracked, but no span is ever written" + recs.hasNext() + recs.next().value() == "local-root-message" + !recs.hasNext() + TEST_DATA_STREAMS_WRITER.waitForGroups(2) + TEST_WRITER.isEmpty() + + cleanup: + consumer?.close() + producer?.close() + } +} + +/** + * Regression guard: producing from inside an already active local trace in DSM-only mode must not + * create any additional span either - the surrounding customer trace is untouched. + */ +class KafkaClientDataStreamsOnlyActiveLocalTraceForkedTest extends KafkaClientDataStreamsOnlyForkedTest { + @Override + void configurePreAgent() { + super.configurePreAgent() + injectSysConfig("integrations.enabled", "false") + } + + def "producing inside an active local trace adds no span to that trace"() { + setup: + def producer = newProducer() + + when: "a message is produced from within an already active local trace" + runUnderTrace("parent") { + producer.send(new ProducerRecord(SHARED_TOPIC, 0, null, "in-active-trace")).get() + } + + then: "the surrounding customer trace contains only its own span, no kafka.produce span" + TEST_WRITER.waitForTraces(1) + TEST_WRITER[0].size() == 1 + TEST_WRITER[0][0].operationName.toString() == "parent" + + cleanup: + producer?.close() + } +} + +/** + * Regression guard for the poll-span suppression site: KafkaConsumerInfoInstrumentation's + * RecordsAdvice used to create a standalone "kafka.poll" span/trace around every consumer.poll() + * call whenever DSM is enabled, regardless of whether Kafka APM tracing itself was enabled. In + * DSM-only mode that span must no longer be created at all. + */ +class KafkaClientDataStreamsOnlyPollSpanForkedTest extends KafkaClientDataStreamsOnlyForkedTest { + static final ListWriter.Filter ACCEPT_ALL = new ListWriter.Filter() { + @Override + boolean accept(List trace) { + return true + } + } + + @Override + void configurePreAgent() { + super.configurePreAgent() + injectSysConfig("integrations.enabled", "false") + } + + def "no poll span is created when kafka tracing is disabled and DSM is enabled"() { + setup: + // Undo any default poll-trace filter: this test needs to see the "kafka.poll" trace, if one + // were (incorrectly) written. + TEST_WRITER.setFilter(ACCEPT_ALL) + def consumerProperties = KafkaTestUtils.consumerProps("sender", "false", embeddedKafka) + consumerProperties.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest") + def consumer = new KafkaConsumer(consumerProperties) + consumer.assign(Arrays.asList(new TopicPartition(SHARED_TOPIC, 0))) + + when: "the consumer polls with no active local trace and no records to consume" + KafkaTestUtils.getRecords(consumer) + + then: "no standalone kafka.poll trace was written" + TEST_WRITER.isEmpty() + + cleanup: + consumer?.close() + } +} diff --git a/dd-java-agent/instrumentation/kafka/kafka-clients-0.11/src/test/groovy/KafkaClientTestBase.groovy b/dd-java-agent/instrumentation/kafka/kafka-clients-0.11/src/test/groovy/KafkaClientTestBase.groovy index 6d311a1b8d5..17c3e665f67 100644 --- a/dd-java-agent/instrumentation/kafka/kafka-clients-0.11/src/test/groovy/KafkaClientTestBase.groovy +++ b/dd-java-agent/instrumentation/kafka/kafka-clients-0.11/src/test/groovy/KafkaClientTestBase.groovy @@ -1544,3 +1544,4 @@ class KafkaClientBadBase64HeaderForkedTest extends InstrumentationSpecification producer?.close() } } + diff --git a/dd-java-agent/instrumentation/kafka/kafka-clients-3.8/src/main/java/datadog/trace/instrumentation/kafka_clients38/ConsumerCoordinatorInstrumentation.java b/dd-java-agent/instrumentation/kafka/kafka-clients-3.8/src/main/java/datadog/trace/instrumentation/kafka_clients38/ConsumerCoordinatorInstrumentation.java index f2d99348473..0df853718bf 100644 --- a/dd-java-agent/instrumentation/kafka/kafka-clients-3.8/src/main/java/datadog/trace/instrumentation/kafka_clients38/ConsumerCoordinatorInstrumentation.java +++ b/dd-java-agent/instrumentation/kafka/kafka-clients-3.8/src/main/java/datadog/trace/instrumentation/kafka_clients38/ConsumerCoordinatorInstrumentation.java @@ -13,11 +13,11 @@ import net.bytebuddy.matcher.ElementMatcher; @AutoService(InstrumenterModule.class) -public final class ConsumerCoordinatorInstrumentation extends InstrumenterModule.Tracing +public final class ConsumerCoordinatorInstrumentation extends InstrumenterModule.DataStreams implements Instrumenter.ForSingleType, Instrumenter.HasMethodAdvice { public ConsumerCoordinatorInstrumentation() { - super("kafka", "kafka-3.8"); + super(KafkaDecorator.INTEGRATION_NAME, KafkaDecorator.LEGACY_INTEGRATION_NAME); } @Override diff --git a/dd-java-agent/instrumentation/kafka/kafka-clients-3.8/src/main/java/datadog/trace/instrumentation/kafka_clients38/KafkaConsumerInfoInstrumentation.java b/dd-java-agent/instrumentation/kafka/kafka-clients-3.8/src/main/java/datadog/trace/instrumentation/kafka_clients38/KafkaConsumerInfoInstrumentation.java index c94c49369ed..e824b1af9ae 100644 --- a/dd-java-agent/instrumentation/kafka/kafka-clients-3.8/src/main/java/datadog/trace/instrumentation/kafka_clients38/KafkaConsumerInfoInstrumentation.java +++ b/dd-java-agent/instrumentation/kafka/kafka-clients-3.8/src/main/java/datadog/trace/instrumentation/kafka_clients38/KafkaConsumerInfoInstrumentation.java @@ -24,13 +24,13 @@ * and cluster ID, in the context store for later use. */ @AutoService(InstrumenterModule.class) -public final class KafkaConsumerInfoInstrumentation extends InstrumenterModule.Tracing +public final class KafkaConsumerInfoInstrumentation extends InstrumenterModule.DataStreams implements Instrumenter.ForTypeHierarchy, Instrumenter.HasMethodAdvice, Instrumenter.WithTypeStructure { public KafkaConsumerInfoInstrumentation() { - super("kafka", "kafka-3.8"); + super(KafkaDecorator.INTEGRATION_NAME, KafkaDecorator.LEGACY_INTEGRATION_NAME); } @Override diff --git a/dd-java-agent/instrumentation/kafka/kafka-clients-3.8/src/main/java/datadog/trace/instrumentation/kafka_clients38/KafkaConsumerInstrumentation.java b/dd-java-agent/instrumentation/kafka/kafka-clients-3.8/src/main/java/datadog/trace/instrumentation/kafka_clients38/KafkaConsumerInstrumentation.java index 70f3f6dbc92..202fcd37770 100644 --- a/dd-java-agent/instrumentation/kafka/kafka-clients-3.8/src/main/java/datadog/trace/instrumentation/kafka_clients38/KafkaConsumerInstrumentation.java +++ b/dd-java-agent/instrumentation/kafka/kafka-clients-3.8/src/main/java/datadog/trace/instrumentation/kafka_clients38/KafkaConsumerInstrumentation.java @@ -19,11 +19,11 @@ import net.bytebuddy.matcher.ElementMatcher; @AutoService(InstrumenterModule.class) -public final class KafkaConsumerInstrumentation extends InstrumenterModule.Tracing +public final class KafkaConsumerInstrumentation extends InstrumenterModule.DataStreams implements Instrumenter.ForSingleType, Instrumenter.HasMethodAdvice { public KafkaConsumerInstrumentation() { - super("kafka", "kafka-3.8"); + super(KafkaDecorator.INTEGRATION_NAME, KafkaDecorator.LEGACY_INTEGRATION_NAME); } @Override diff --git a/dd-java-agent/instrumentation/kafka/kafka-clients-3.8/src/main/java/datadog/trace/instrumentation/kafka_clients38/KafkaProducerInstrumentation.java b/dd-java-agent/instrumentation/kafka/kafka-clients-3.8/src/main/java/datadog/trace/instrumentation/kafka_clients38/KafkaProducerInstrumentation.java index 0680c757e37..a5a25128f3a 100644 --- a/dd-java-agent/instrumentation/kafka/kafka-clients-3.8/src/main/java/datadog/trace/instrumentation/kafka_clients38/KafkaProducerInstrumentation.java +++ b/dd-java-agent/instrumentation/kafka/kafka-clients-3.8/src/main/java/datadog/trace/instrumentation/kafka_clients38/KafkaProducerInstrumentation.java @@ -16,11 +16,11 @@ import net.bytebuddy.matcher.ElementMatcher; @AutoService(InstrumenterModule.class) -public final class KafkaProducerInstrumentation extends InstrumenterModule.Tracing +public final class KafkaProducerInstrumentation extends InstrumenterModule.DataStreams implements Instrumenter.ForSingleType, Instrumenter.HasMethodAdvice { public KafkaProducerInstrumentation() { - super("kafka", "kafka-3.8"); + super(KafkaDecorator.INTEGRATION_NAME, KafkaDecorator.LEGACY_INTEGRATION_NAME); } @Override diff --git a/dd-java-agent/instrumentation/kafka/kafka-clients-3.8/src/main/java/datadog/trace/instrumentation/kafka_clients38/LegacyKafkaConsumerInfoInstrumentation.java b/dd-java-agent/instrumentation/kafka/kafka-clients-3.8/src/main/java/datadog/trace/instrumentation/kafka_clients38/LegacyKafkaConsumerInfoInstrumentation.java index dd36ff1d934..67833947dab 100644 --- a/dd-java-agent/instrumentation/kafka/kafka-clients-3.8/src/main/java/datadog/trace/instrumentation/kafka_clients38/LegacyKafkaConsumerInfoInstrumentation.java +++ b/dd-java-agent/instrumentation/kafka/kafka-clients-3.8/src/main/java/datadog/trace/instrumentation/kafka_clients38/LegacyKafkaConsumerInfoInstrumentation.java @@ -24,13 +24,13 @@ * and cluster ID, in the context store for later use. */ @AutoService(InstrumenterModule.class) -public final class LegacyKafkaConsumerInfoInstrumentation extends InstrumenterModule.Tracing +public final class LegacyKafkaConsumerInfoInstrumentation extends InstrumenterModule.DataStreams implements Instrumenter.ForTypeHierarchy, Instrumenter.HasMethodAdvice, Instrumenter.WithTypeStructure { public LegacyKafkaConsumerInfoInstrumentation() { - super("kafka", "kafka-3.8"); + super(KafkaDecorator.INTEGRATION_NAME, KafkaDecorator.LEGACY_INTEGRATION_NAME); } @Override diff --git a/dd-java-agent/instrumentation/kafka/kafka-clients-3.8/src/main/java/datadog/trace/instrumentation/kafka_clients38/MessageListenerInstrumentation.java b/dd-java-agent/instrumentation/kafka/kafka-clients-3.8/src/main/java/datadog/trace/instrumentation/kafka_clients38/MessageListenerInstrumentation.java index da17835a1a2..9da06603724 100644 --- a/dd-java-agent/instrumentation/kafka/kafka-clients-3.8/src/main/java/datadog/trace/instrumentation/kafka_clients38/MessageListenerInstrumentation.java +++ b/dd-java-agent/instrumentation/kafka/kafka-clients-3.8/src/main/java/datadog/trace/instrumentation/kafka_clients38/MessageListenerInstrumentation.java @@ -12,12 +12,19 @@ import net.bytebuddy.description.type.TypeDescription; import net.bytebuddy.matcher.ElementMatcher; +/** + * Applies Code Origin (span origin) advice to Spring Kafka message listeners. This is a purely-APM + * concern with no Data Streams behaviour, so it deliberately stays on the {@link + * InstrumenterModule.Tracing} base class: {@link InstrumenterModule.DataStreams} ORs {@code + * isDataStreamsEnabled()} into {@code isEnabled()}, which would install this instrumentation in + * DSM-only deployments where Kafka tracing is off. + */ @AutoService(InstrumenterModule.class) public class MessageListenerInstrumentation extends InstrumenterModule.Tracing implements Instrumenter.ForTypeHierarchy, Instrumenter.HasMethodAdvice { public MessageListenerInstrumentation() { - super("kafka", "kafka-3.8"); + super(KafkaDecorator.INTEGRATION_NAME, KafkaDecorator.LEGACY_INTEGRATION_NAME); } @Override diff --git a/dd-java-agent/instrumentation/kafka/kafka-clients-3.8/src/main/java/datadog/trace/instrumentation/kafka_clients38/MetadataInstrumentation.java b/dd-java-agent/instrumentation/kafka/kafka-clients-3.8/src/main/java/datadog/trace/instrumentation/kafka_clients38/MetadataInstrumentation.java index 3907ad0c18c..7bab5497904 100644 --- a/dd-java-agent/instrumentation/kafka/kafka-clients-3.8/src/main/java/datadog/trace/instrumentation/kafka_clients38/MetadataInstrumentation.java +++ b/dd-java-agent/instrumentation/kafka/kafka-clients-3.8/src/main/java/datadog/trace/instrumentation/kafka_clients38/MetadataInstrumentation.java @@ -15,11 +15,11 @@ import net.bytebuddy.matcher.ElementMatcher; @AutoService(InstrumenterModule.class) -public class MetadataInstrumentation extends InstrumenterModule.Tracing +public class MetadataInstrumentation extends InstrumenterModule.DataStreams implements Instrumenter.ForTypeHierarchy, Instrumenter.HasMethodAdvice { public MetadataInstrumentation() { - super("kafka", "kafka-3.8"); + super(KafkaDecorator.INTEGRATION_NAME, KafkaDecorator.LEGACY_INTEGRATION_NAME); } @Override diff --git a/dd-java-agent/instrumentation/kafka/kafka-clients-3.8/src/main/java/datadog/trace/instrumentation/kafka_clients38/OffsetCommitCallbackInvokerInstrumentation.java b/dd-java-agent/instrumentation/kafka/kafka-clients-3.8/src/main/java/datadog/trace/instrumentation/kafka_clients38/OffsetCommitCallbackInvokerInstrumentation.java index e64beeb1ebd..62dd9416b85 100644 --- a/dd-java-agent/instrumentation/kafka/kafka-clients-3.8/src/main/java/datadog/trace/instrumentation/kafka_clients38/OffsetCommitCallbackInvokerInstrumentation.java +++ b/dd-java-agent/instrumentation/kafka/kafka-clients-3.8/src/main/java/datadog/trace/instrumentation/kafka_clients38/OffsetCommitCallbackInvokerInstrumentation.java @@ -10,10 +10,10 @@ // new - this instrumentation is completely new. // the purpose of this class is to provide us with information on consumer group and cluster ID -public class OffsetCommitCallbackInvokerInstrumentation extends InstrumenterModule.Tracing +public class OffsetCommitCallbackInvokerInstrumentation extends InstrumenterModule.DataStreams implements Instrumenter.ForSingleType, Instrumenter.HasMethodAdvice { public OffsetCommitCallbackInvokerInstrumentation() { - super("kafka", "kafka-3.8"); + super(KafkaDecorator.INTEGRATION_NAME, KafkaDecorator.LEGACY_INTEGRATION_NAME); } @Override diff --git a/dd-java-agent/instrumentation/kafka/kafka-clients-3.8/src/main/java17/datadog/trace/instrumentation/kafka_clients38/KafkaDecorator.java b/dd-java-agent/instrumentation/kafka/kafka-clients-3.8/src/main/java17/datadog/trace/instrumentation/kafka_clients38/KafkaDecorator.java index d2d6f53b8a9..1a1172d264d 100644 --- a/dd-java-agent/instrumentation/kafka/kafka-clients-3.8/src/main/java17/datadog/trace/instrumentation/kafka_clients38/KafkaDecorator.java +++ b/dd-java-agent/instrumentation/kafka/kafka-clients-3.8/src/main/java17/datadog/trace/instrumentation/kafka_clients38/KafkaDecorator.java @@ -12,6 +12,7 @@ import datadog.trace.api.Config; import datadog.trace.api.Functions; +import datadog.trace.api.InstrumenterConfig; import datadog.trace.api.cache.DDCache; import datadog.trace.api.cache.DDCaches; import datadog.trace.api.naming.SpanNaming; @@ -20,6 +21,7 @@ import datadog.trace.bootstrap.instrumentation.api.Tags; import datadog.trace.bootstrap.instrumentation.api.UTF8BytesString; import datadog.trace.bootstrap.instrumentation.decorator.MessagingClientDecorator; +import java.util.Arrays; import java.util.function.Function; import java.util.function.Supplier; import org.apache.kafka.clients.consumer.ConsumerRecord; @@ -29,6 +31,10 @@ public class KafkaDecorator extends MessagingClientDecorator { private static final String KAFKA = "kafka"; + // Kept in sync with the names each kafka-clients-3.8 instrumentation module passes to its own + // super(...) constructor call, so TRACING_ENABLED can't drift from what is actually registered. + public static final String INTEGRATION_NAME = KAFKA; + public static final String LEGACY_INTEGRATION_NAME = "kafka-3.8"; public static final CharSequence JAVA_KAFKA = UTF8BytesString.create("java-kafka"); public static final CharSequence KAFKA_CONSUME = UTF8BytesString.create( @@ -42,6 +48,11 @@ public class KafkaDecorator extends MessagingClientDecorator { public static final boolean KAFKA_LEGACY_TRACING = Config.get().isKafkaLegacyTracingEnabled(); public static final boolean TIME_IN_QUEUE_ENABLED = Config.get().isTimeInQueueEnabled(!KAFKA_LEGACY_TRACING, KAFKA); + public static final boolean TRACING_ENABLED = + InstrumenterConfig.get() + .isIntegrationEnabled( + Arrays.asList(INTEGRATION_NAME, LEGACY_INTEGRATION_NAME), + InstrumenterConfig.get().isIntegrationsEnabled()); public static final String KAFKA_PRODUCED_KEY = "x_datadog_kafka_produced"; private final String spanKind; private final CharSequence spanType; diff --git a/dd-java-agent/instrumentation/kafka/kafka-clients-3.8/src/main/java17/datadog/trace/instrumentation/kafka_clients38/ProducerAdvice.java b/dd-java-agent/instrumentation/kafka/kafka-clients-3.8/src/main/java17/datadog/trace/instrumentation/kafka_clients38/ProducerAdvice.java index 01905ee65e4..c2b2fc4725d 100644 --- a/dd-java-agent/instrumentation/kafka/kafka-clients-3.8/src/main/java17/datadog/trace/instrumentation/kafka_clients38/ProducerAdvice.java +++ b/dd-java-agent/instrumentation/kafka/kafka-clients-3.8/src/main/java17/datadog/trace/instrumentation/kafka_clients38/ProducerAdvice.java @@ -4,6 +4,7 @@ import static datadog.trace.bootstrap.instrumentation.api.AgentTracer.activateSpan; import static datadog.trace.bootstrap.instrumentation.api.AgentTracer.activeSpan; import static datadog.trace.bootstrap.instrumentation.api.AgentTracer.startSpan; +import static datadog.trace.bootstrap.instrumentation.api.AgentTracer.traceConfig; import static datadog.trace.instrumentation.kafka_clients38.KafkaDecorator.JAVA_KAFKA; import static datadog.trace.instrumentation.kafka_clients38.KafkaDecorator.KAFKA_PRODUCE; import static datadog.trace.instrumentation.kafka_clients38.KafkaDecorator.PRODUCER_DECORATE; @@ -15,6 +16,7 @@ import datadog.trace.bootstrap.instrumentation.api.InstrumentationTags; import datadog.trace.instrumentation.kafka_common.ClusterIdHolder; import datadog.trace.instrumentation.kafka_common.MetadataState; +import datadog.trace.instrumentation.kafka_common.Utils; import net.bytebuddy.asm.Advice; import org.apache.kafka.clients.Metadata; import org.apache.kafka.clients.producer.Callback; @@ -44,20 +46,32 @@ public static AgentScope onEnter( final AgentSpanContext extractedContext = extractContextAndGetSpanContext(record.headers(), TextMapExtractAdapter.GETTER); - final AgentSpan localActiveSpan = activeSpan(); - final AgentSpan span; final AgentSpan callbackParentSpan; - if (extractedContext != null) { - span = startSpan(JAVA_KAFKA.toString(), KAFKA_PRODUCE, extractedContext); + if (!KafkaDecorator.TRACING_ENABLED && traceConfig().isDataStreamsEnabled()) { + // DSM-only mode: never create a real span, so nothing for this integration is ever + // written to the agent. The pathway is carried on a lightweight, never-collected span + // shim instead, falling back to whatever pathway the currently active span (if any) + // is carrying so a consume->produce chain keeps propagating the same pathway. + final AgentSpan localActiveSpan = activeSpan(); + final AgentSpanContext pathwaySource = + extractedContext != null + ? extractedContext + : localActiveSpan == null ? null : localActiveSpan.spanContext(); + span = Utils.newPathwayOnlySpan(pathwaySource); callbackParentSpan = span; } else { - span = startSpan(JAVA_KAFKA.toString(), KAFKA_PRODUCE); - callbackParentSpan = localActiveSpan; + if (extractedContext != null) { + span = startSpan(JAVA_KAFKA.toString(), KAFKA_PRODUCE, extractedContext); + callbackParentSpan = span; + } else { + span = startSpan(JAVA_KAFKA.toString(), KAFKA_PRODUCE); + callbackParentSpan = activeSpan(); + } + PRODUCER_DECORATE.afterStart(span); + PRODUCER_DECORATE.onProduce(span, record, producerConfig, clusterId); } - PRODUCER_DECORATE.afterStart(span); - PRODUCER_DECORATE.onProduce(span, record, producerConfig, clusterId); callback = new KafkaProducerCallback(callback, callbackParentSpan, span, clusterId); diff --git a/dd-java-agent/instrumentation/kafka/kafka-clients-3.8/src/main/java17/datadog/trace/instrumentation/kafka_clients38/RecordsAdvice.java b/dd-java-agent/instrumentation/kafka/kafka-clients-3.8/src/main/java17/datadog/trace/instrumentation/kafka_clients38/RecordsAdvice.java index b8f3dff049a..2884eaecf9d 100644 --- a/dd-java-agent/instrumentation/kafka/kafka-clients-3.8/src/main/java17/datadog/trace/instrumentation/kafka_clients38/RecordsAdvice.java +++ b/dd-java-agent/instrumentation/kafka/kafka-clients-3.8/src/main/java17/datadog/trace/instrumentation/kafka_clients38/RecordsAdvice.java @@ -39,7 +39,10 @@ public static AgentScope onEnter(@Advice.This ConsumerDelegate consumer) { } } - if (traceConfig().isDataStreamsEnabled()) { + if (traceConfig().isDataStreamsEnabled() && KafkaDecorator.TRACING_ENABLED) { + // DSM-only mode (tracing disabled) never creates a real poll span: TracingIterator carries + // its pathway context on a lightweight, never-collected span shim instead, so there's + // nothing here that needs wrapping/protecting from being force-dropped. final AgentSpan span = startSpan(JAVA_KAFKA.toString(), KAFKA_POLL); return activateSpan(span); } diff --git a/dd-java-agent/instrumentation/kafka/kafka-clients-3.8/src/main/java17/datadog/trace/instrumentation/kafka_clients38/TracingIterator.java b/dd-java-agent/instrumentation/kafka/kafka-clients-3.8/src/main/java17/datadog/trace/instrumentation/kafka_clients38/TracingIterator.java index 7beba848473..524accbe86f 100644 --- a/dd-java-agent/instrumentation/kafka/kafka-clients-3.8/src/main/java17/datadog/trace/instrumentation/kafka_clients38/TracingIterator.java +++ b/dd-java-agent/instrumentation/kafka/kafka-clients-3.8/src/main/java17/datadog/trace/instrumentation/kafka_clients38/TracingIterator.java @@ -93,56 +93,11 @@ protected void startNewRecordSpan(ConsumerRecord val) { previousSpan.finishWithEndToEnd(); } } - AgentSpan span, queueSpan = null; if (val != null) { - if (!Config.get().isKafkaClientPropagationDisabledForTopic(val.topic())) { - final AgentSpanContext spanContext = - extractContextAndGetSpanContext(val.headers(), GETTER); - long timeInQueueStart = GETTER.extractTimeInQueueStart(val.headers()); - if (timeInQueueStart == 0 || !KafkaDecorator.TIME_IN_QUEUE_ENABLED) { - span = startSpan(JAVA_KAFKA.toString(), operationName, spanContext); - } else { - queueSpan = - startSpan( - JAVA_KAFKA.toString(), - KafkaDecorator.KAFKA_DELIVER, - spanContext, - MILLISECONDS.toMicros(timeInQueueStart)); - KafkaDecorator.BROKER_DECORATE.afterStart(queueSpan); - KafkaDecorator.BROKER_DECORATE.onTimeInQueue(queueSpan, val); - span = startSpan(JAVA_KAFKA.toString(), operationName, queueSpan.spanContext()); - KafkaDecorator.BROKER_DECORATE.beforeFinish(queueSpan); - // The queueSpan will be finished after inner span has been activated to ensure that - // spans are written out together by TraceStructureWriter when running in strict mode - } - - DataStreamsTags tags = create("kafka", INBOUND, val.topic(), group, clusterId); - final long payloadSize = - traceConfig().isDataStreamsEnabled() ? Utils.computePayloadSizeBytes(val) : 0; - if (StreamingContext.STREAMING_CONTEXT.isDisabledForTopic(val.topic())) { - AgentTracer.get() - .getDataStreamsMonitoring() - .setCheckpoint(span, create(tags, val.timestamp(), payloadSize)); - } else { - // when we're in a streaming context we want to consume only from source topics - if (StreamingContext.STREAMING_CONTEXT.isSourceTopic(val.topic())) { - // We have to inject the context to headers here, - // since the data received from the source may leave the topology on - // some other instance of the application, breaking the context propagation - // for DSM users - Propagator dsmPropagator = Propagators.forConcern(DSM_CONCERN); - DataStreamsContext dsmContext = create(tags, val.timestamp(), payloadSize); - dsmPropagator.inject(span.with(dsmContext), val.headers(), SETTER); - } - } - } else { - span = startSpan(JAVA_KAFKA.toString(), operationName, null); - } - if (val.value() == null) { - span.setTag(InstrumentationTags.TOMBSTONE, true); - } - decorator.afterStart(span); - decorator.onConsume(span, val, group, clusterId, bootstrapServers); + final AgentSpan span = + !KafkaDecorator.TRACING_ENABLED && traceConfig().isDataStreamsEnabled() + ? startDsmOnlyPathwaySpan(val) + : startTracedConsumeSpan(val); if (InstrumenterConfig.get().isLegacyContextManagerEnabled()) { activateNext(span); } else { @@ -151,23 +106,108 @@ protected void startNewRecordSpan(ConsumerRecord val) { previousSpan.finishWithEndToEnd(); } } - if (null != queueSpan) { - queueSpan.finish(); - } - - AgentTracer.get() - .getDataStreamsMonitoring() - .trackTransaction( - span, - DataStreamsTransactionExtractor.Type.KAFKA_CONSUME_HEADERS, - val.headers(), - Utils.DSM_TRANSACTION_SOURCE_READER); } } catch (final Exception e) { log.debug("Error starting new record span", e); } } + /** + * Creates and activates the real APM consume span (and, when time-in-queue is enabled, its broker + * parent), tags it, and reports DSM checkpoints/transactions off of it. + */ + private AgentSpan startTracedConsumeSpan(ConsumerRecord val) { + AgentSpan span, queueSpan = null; + if (!Config.get().isKafkaClientPropagationDisabledForTopic(val.topic())) { + final AgentSpanContext spanContext = extractContextAndGetSpanContext(val.headers(), GETTER); + long timeInQueueStart = GETTER.extractTimeInQueueStart(val.headers()); + if (timeInQueueStart == 0 || !KafkaDecorator.TIME_IN_QUEUE_ENABLED) { + span = startSpan(JAVA_KAFKA.toString(), operationName, spanContext); + } else { + queueSpan = + startSpan( + JAVA_KAFKA.toString(), + KafkaDecorator.KAFKA_DELIVER, + spanContext, + MILLISECONDS.toMicros(timeInQueueStart)); + KafkaDecorator.BROKER_DECORATE.afterStart(queueSpan); + KafkaDecorator.BROKER_DECORATE.onTimeInQueue(queueSpan, val); + span = startSpan(JAVA_KAFKA.toString(), operationName, queueSpan.spanContext()); + KafkaDecorator.BROKER_DECORATE.beforeFinish(queueSpan); + // The queueSpan will be finished after inner span has been activated to ensure that + // spans are written out together by TraceStructureWriter when running in strict mode + } + + DataStreamsTags tags = create("kafka", INBOUND, val.topic(), group, clusterId); + final long payloadSize = + traceConfig().isDataStreamsEnabled() ? Utils.computePayloadSizeBytes(val) : 0; + reportDsmCheckpointOrInject(span, val, tags, payloadSize); + } else { + span = startSpan(JAVA_KAFKA.toString(), operationName, null); + } + if (val.value() == null) { + span.setTag(InstrumentationTags.TOMBSTONE, true); + } + decorator.afterStart(span); + decorator.onConsume(span, val, group, clusterId, bootstrapServers); + if (null != queueSpan) { + queueSpan.finish(); + } + + trackDsmConsumeTransaction(span, val); + return span; + } + + /** + * DSM-only mode (tracing disabled for kafka, DSM enabled): never creates a real span, so no span + * is ever written to the agent for this integration. Only the pathway checkpoint/injection and + * transaction tracking happen, carried by a lightweight, never-collected span shim. + */ + private AgentSpan startDsmOnlyPathwaySpan(ConsumerRecord val) { + AgentSpan span; + if (!Config.get().isKafkaClientPropagationDisabledForTopic(val.topic())) { + final AgentSpanContext extractedContext = + extractContextAndGetSpanContext(val.headers(), GETTER); + span = Utils.newPathwayOnlySpan(extractedContext); + DataStreamsTags tags = create("kafka", INBOUND, val.topic(), group, clusterId); + final long payloadSize = Utils.computePayloadSizeBytes(val); + reportDsmCheckpointOrInject(span, val, tags, payloadSize); + } else { + span = Utils.newPathwayOnlySpan(null); + } + trackDsmConsumeTransaction(span, val); + return span; + } + + /** + * Reports a DSM checkpoint for {@code val}'s topic, or - when in a streaming context and {@code + * val}'s topic is a source topic - injects the pathway context into its headers so it survives + * leaving the topology on another instance of the application. + */ + private void reportDsmCheckpointOrInject( + AgentSpan span, ConsumerRecord val, DataStreamsTags tags, long payloadSize) { + if (StreamingContext.STREAMING_CONTEXT.isDisabledForTopic(val.topic())) { + AgentTracer.get() + .getDataStreamsMonitoring() + .setCheckpoint(span, create(tags, val.timestamp(), payloadSize)); + } else if (StreamingContext.STREAMING_CONTEXT.isSourceTopic(val.topic())) { + // when we're in a streaming context we want to consume only from source topics + Propagator dsmPropagator = Propagators.forConcern(DSM_CONCERN); + DataStreamsContext dsmContext = create(tags, val.timestamp(), payloadSize); + dsmPropagator.inject(span.with(dsmContext), val.headers(), SETTER); + } + } + + private void trackDsmConsumeTransaction(AgentSpan span, ConsumerRecord val) { + AgentTracer.get() + .getDataStreamsMonitoring() + .trackTransaction( + span, + DataStreamsTransactionExtractor.Type.KAFKA_CONSUME_HEADERS, + val.headers(), + Utils.DSM_TRANSACTION_SOURCE_READER); + } + @Override public void remove() { delegateIterator.remove(); diff --git a/dd-java-agent/instrumentation/kafka/kafka-clients-3.8/src/test/groovy/KafkaClientDataStreamsOnlyForkedTest.groovy b/dd-java-agent/instrumentation/kafka/kafka-clients-3.8/src/test/groovy/KafkaClientDataStreamsOnlyForkedTest.groovy new file mode 100644 index 00000000000..0bf9b64a94a --- /dev/null +++ b/dd-java-agent/instrumentation/kafka/kafka-clients-3.8/src/test/groovy/KafkaClientDataStreamsOnlyForkedTest.groovy @@ -0,0 +1,144 @@ +import datadog.trace.agent.test.InstrumentationSpecification +import org.apache.kafka.clients.consumer.ConsumerConfig +import org.apache.kafka.clients.consumer.KafkaConsumer +import org.apache.kafka.clients.producer.KafkaProducer +import org.apache.kafka.clients.producer.ProducerRecord +import org.apache.kafka.common.TopicPartition +import org.apache.kafka.common.serialization.StringSerializer +import org.springframework.kafka.test.EmbeddedKafkaBroker +import org.springframework.kafka.test.EmbeddedKafkaKraftBroker +import org.springframework.kafka.test.utils.KafkaTestUtils + +import static datadog.trace.agent.test.utils.TraceUtils.runUnderTrace + +/** + * DSM-only coverage for kafka-clients-3.8: when Kafka APM tracing is disabled + * (integrations.enabled=false) but DSM is enabled, this integration must never create or write a + * real span/trace, while pathway checkpoints must still be tracked. These specs deliberately do + * not extend {@code KafkaClientTestBase}: that base asserts on APM spans, which are correctly + * absent once Kafka tracing is off. + */ +abstract class KafkaClientDataStreamsOnlyForkedTest extends InstrumentationSpecification { + static final SHARED_TOPIC = "shared.topic" + + EmbeddedKafkaBroker embeddedKafka + + def setup() { + embeddedKafka = new EmbeddedKafkaKraftBroker(1, 2, SHARED_TOPIC) + embeddedKafka.afterPropertiesSet() + } + + def cleanup() { + embeddedKafka.destroy() + } + + @Override + boolean useStrictTraceWrites() { + return false + } + + @Override + protected boolean isDataStreamsEnabled() { + return true + } + + @Override + void configurePreAgent() { + super.configurePreAgent() + injectSysConfig("integrations.enabled", "false") + } + + protected KafkaProducer newProducer() { + return new KafkaProducer( + KafkaTestUtils.producerProps(embeddedKafka.getBrokersAsString()), + new StringSerializer(), + new StringSerializer()) + } +} + +/** + * Regression guard for the contract that "integration disabled => zero spans of that type ever + * reach the agent": producing and consuming a message in DSM-only mode must not write any trace, + * even though DSM checkpoints for the same produce/consume are still tracked. + */ +class KafkaClientDataStreamsOnlyNoSpansForkedTest extends KafkaClientDataStreamsOnlyForkedTest { + + def "produce and consume in DSM-only mode write no spans, but still track DSM checkpoints"() { + setup: + def kafkaPartition = 0 + def consumerProperties = KafkaTestUtils.consumerProps("sender", "false", embeddedKafka) + consumerProperties.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest") + def consumer = new KafkaConsumer(consumerProperties) + def producer = newProducer() + consumer.assign(Arrays.asList(new TopicPartition(SHARED_TOPIC, kafkaPartition))) + + when: "a message is produced and consumed with no active trace" + producer.send(new ProducerRecord(SHARED_TOPIC, kafkaPartition, null, "dsm-only-message")).get() + def recs = KafkaTestUtils.getRecords(consumer) + .records(new TopicPartition(SHARED_TOPIC, kafkaPartition)).iterator() + + then: "the message is delivered and DSM checkpoints are tracked for both hops" + recs.hasNext() + recs.next().value() == "dsm-only-message" + !recs.hasNext() + TEST_DATA_STREAMS_WRITER.waitForGroups(2, 15000) + + and: "no span was ever created or written for this integration" + TEST_WRITER.isEmpty() + + cleanup: + consumer?.close() + producer?.close() + } +} + +/** + * Regression guard: producing from inside an already active local trace in DSM-only mode must not + * create any additional span either - the surrounding customer trace is untouched. + */ +class KafkaClientDataStreamsOnlyActiveLocalTraceForkedTest extends KafkaClientDataStreamsOnlyForkedTest { + + def "producing inside an active local trace in DSM-only mode adds no span to that trace"() { + setup: + def producer = newProducer() + + when: "a message is produced from within an already active local trace" + runUnderTrace("parent") { + producer.send(new ProducerRecord(SHARED_TOPIC, 0, null, "in-active-trace")).get() + } + + then: "the surrounding customer trace contains only its own span, no kafka.produce span" + TEST_WRITER.waitForTraces(1) + TEST_WRITER[0].size() == 1 + TEST_WRITER[0][0].operationName.toString() == "parent" + + cleanup: + producer?.close() + } +} + +/** + * Regression guard for the poll-span suppression site: KafkaConsumerInfoInstrumentation's + * RecordsAdvice used to create a standalone "kafka.poll" span/trace around every consumer.poll() + * call whenever DSM is enabled, regardless of whether Kafka APM tracing itself was enabled. In + * DSM-only mode that span must no longer be created at all. + */ +class KafkaClientDataStreamsOnlyPollSpanForkedTest extends KafkaClientDataStreamsOnlyForkedTest { + + def "no poll span is created in DSM-only mode"() { + setup: + def consumerProperties = KafkaTestUtils.consumerProps("sender", "false", embeddedKafka) + consumerProperties.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest") + def consumer = new KafkaConsumer(consumerProperties) + consumer.assign(Arrays.asList(new TopicPartition(SHARED_TOPIC, 0))) + + when: "the consumer polls with no active local trace and no records to consume" + KafkaTestUtils.getRecords(consumer) + + then: "no standalone kafka.poll trace was written" + TEST_WRITER.isEmpty() + + cleanup: + consumer?.close() + } +} diff --git a/dd-java-agent/instrumentation/kafka/kafka-common/src/main/java/datadog/trace/instrumentation/kafka_common/Utils.java b/dd-java-agent/instrumentation/kafka/kafka-common/src/main/java/datadog/trace/instrumentation/kafka_common/Utils.java index d0f7eb4fcea..51648f299af 100644 --- a/dd-java-agent/instrumentation/kafka/kafka-common/src/main/java/datadog/trace/instrumentation/kafka_common/Utils.java +++ b/dd-java-agent/instrumentation/kafka/kafka-common/src/main/java/datadog/trace/instrumentation/kafka_common/Utils.java @@ -1,6 +1,11 @@ package datadog.trace.instrumentation.kafka_common; import datadog.trace.api.datastreams.DataStreamsTransactionTracker; +import datadog.trace.api.datastreams.PathwayContext; +import datadog.trace.bootstrap.instrumentation.api.AgentSpan; +import datadog.trace.bootstrap.instrumentation.api.AgentSpanContext; +import datadog.trace.bootstrap.instrumentation.api.AgentTracer; +import datadog.trace.bootstrap.instrumentation.api.TagContext; import java.nio.charset.StandardCharsets; import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.kafka.common.header.Header; @@ -9,6 +14,21 @@ public final class Utils { private Utils() {} // prevent instantiation + /** + * Builds a span-shaped carrier for a {@link PathwayContext} only, without creating a real trace: + * no sampling, no trace-collector registration, never written to the agent. Used when APM tracing + * is disabled for the integration but DSM is enabled, so a pathway can still be + * propagated/checkpointed through the normal active-span-based DSM APIs. + */ + public static AgentSpan newPathwayOnlySpan(AgentSpanContext extractedContext) { + PathwayContext pathwayContext = + extractedContext == null ? null : extractedContext.getPathwayContext(); + if (pathwayContext == null) { + pathwayContext = AgentTracer.get().getDataStreamsMonitoring().newPathwayContext(); + } + return AgentSpan.fromSpanContext(new TagContext().withPathwayContext(pathwayContext)); + } + public static DataStreamsTransactionTracker.TransactionSourceReader DSM_TRANSACTION_SOURCE_READER = (source, headerName) -> { diff --git a/dd-java-agent/instrumentation/kafka/kafka-connect-0.11/src/main/java/datadog/trace/instrumentation/kafka_connect/ConnectWorkerInstrumentation.java b/dd-java-agent/instrumentation/kafka/kafka-connect-0.11/src/main/java/datadog/trace/instrumentation/kafka_connect/ConnectWorkerInstrumentation.java index 5c56e1341d0..bf2b04c21e6 100644 --- a/dd-java-agent/instrumentation/kafka/kafka-connect-0.11/src/main/java/datadog/trace/instrumentation/kafka_connect/ConnectWorkerInstrumentation.java +++ b/dd-java-agent/instrumentation/kafka/kafka-connect-0.11/src/main/java/datadog/trace/instrumentation/kafka_connect/ConnectWorkerInstrumentation.java @@ -16,7 +16,7 @@ import org.apache.kafka.connect.util.ConnectorTaskId; @AutoService(InstrumenterModule.class) -public final class ConnectWorkerInstrumentation extends InstrumenterModule.Tracing +public final class ConnectWorkerInstrumentation extends InstrumenterModule.DataStreams implements Instrumenter.ForTypeHierarchy, Instrumenter.HasMethodAdvice { static final String TARGET_TYPE = "org.apache.kafka.connect.runtime.WorkerTask"; diff --git a/dd-java-agent/instrumentation/kafka/kafka-streams-0.11/src/main/java/datadog/trace/instrumentation/kafka_streams/KafkaStreamTaskInstrumentation.java b/dd-java-agent/instrumentation/kafka/kafka-streams-0.11/src/main/java/datadog/trace/instrumentation/kafka_streams/KafkaStreamTaskInstrumentation.java index 81ddfd4202f..20598e3a66f 100644 --- a/dd-java-agent/instrumentation/kafka/kafka-streams-0.11/src/main/java/datadog/trace/instrumentation/kafka_streams/KafkaStreamTaskInstrumentation.java +++ b/dd-java-agent/instrumentation/kafka/kafka-streams-0.11/src/main/java/datadog/trace/instrumentation/kafka_streams/KafkaStreamTaskInstrumentation.java @@ -8,11 +8,13 @@ import static datadog.trace.api.datastreams.DataStreamsTags.createWithGroup; import static datadog.trace.bootstrap.instrumentation.api.AgentPropagation.DSM_CONCERN; import static datadog.trace.bootstrap.instrumentation.api.AgentTracer.activateSpan; +import static datadog.trace.bootstrap.instrumentation.api.AgentTracer.activeSpan; import static datadog.trace.bootstrap.instrumentation.api.AgentTracer.startSpan; import static datadog.trace.bootstrap.instrumentation.api.AgentTracer.traceConfig; import static datadog.trace.bootstrap.instrumentation.api.Java8BytecodeBridge.rootContext; import static datadog.trace.instrumentation.kafka_common.StreamingContext.STREAMING_CONTEXT; import static datadog.trace.instrumentation.kafka_common.Utils.computePayloadSizeBytes; +import static datadog.trace.instrumentation.kafka_common.Utils.newPathwayOnlySpan; import static datadog.trace.instrumentation.kafka_streams.KafkaStreamsDecorator.BROKER_DECORATE; import static datadog.trace.instrumentation.kafka_streams.KafkaStreamsDecorator.CONSUMER_DECORATE; import static datadog.trace.instrumentation.kafka_streams.KafkaStreamsDecorator.JAVA_KAFKA; @@ -57,11 +59,11 @@ import org.apache.kafka.streams.processor.internals.StreamTask; @AutoService(InstrumenterModule.class) -public class KafkaStreamTaskInstrumentation extends InstrumenterModule.Tracing +public class KafkaStreamTaskInstrumentation extends InstrumenterModule.DataStreams implements Instrumenter.ForSingleType, Instrumenter.HasMethodAdvice { public KafkaStreamTaskInstrumentation() { - super("kafka", "kafka-streams"); + super(KafkaStreamsDecorator.INTEGRATION_NAME, KafkaStreamsDecorator.LEGACY_INTEGRATION_NAME); } @Override @@ -276,9 +278,33 @@ public static void start( return; } - AgentSpan span, queueSpan = null; StreamTaskContext streamTaskContext = InstrumentationContext.get(StreamTask.class, StreamTaskContext.class).get(task); + String applicationId = + streamTaskContext != null ? streamTaskContext.getApplicationId() : null; + + final AgentSpan span = + !KafkaStreamsDecorator.TRACING_ENABLED && traceConfig().isDataStreamsEnabled() + ? startDsmOnlyPathwaySpan(record, applicationId) + : startTracedConsumeSpan(record, node, applicationId); + + AgentScope agentScope = activateSpan(span); + + if (streamTaskContext == null) { + streamTaskContext = new StreamTaskContext(); + } + streamTaskContext.setAgentScope(agentScope); + InstrumentationContext.get(StreamTask.class, StreamTaskContext.class) + .put(task, streamTaskContext); + } + + /** + * Creates and activates the real APM consume span (and, when time-in-queue is enabled, its + * broker parent), tags it, and reports DSM checkpoints/transactions off of it. + */ + private static AgentSpan startTracedConsumeSpan( + final StampedRecord record, final ProcessorNode node, final String applicationId) { + AgentSpan span, queueSpan = null; long timeInQueueStart = SR_GETTER.extractTimeInQueueStart(record); if (timeInQueueStart == 0 || !TIME_IN_QUEUE_ENABLED) { span = startSpan(JAVA_KAFKA.toString(), KAFKA_CONSUME); @@ -294,39 +320,56 @@ public static void start( // spans are written out together by TraceStructureWriter when running in strict mode } - String applicationId = null; - if (streamTaskContext != null) { - applicationId = streamTaskContext.getApplicationId(); - } DataStreamsTags tags = createWithGroup("kafka", INBOUND, applicationId, record.topic()); final long payloadSize = traceConfig().isDataStreamsEnabled() ? computePayloadSizeBytes(record.value) : 0; - if (STREAMING_CONTEXT.isDisabledForTopic(record.topic())) { - AgentTracer.get() - .getDataStreamsMonitoring() - .setCheckpoint(span, create(tags, record.timestamp, payloadSize)); - } else { - if (STREAMING_CONTEXT.isSourceTopic(record.topic())) { - Propagator dsmPropagator = Propagators.forConcern(DSM_CONCERN); - DataStreamsContext dsmContext = create(tags, record.timestamp, payloadSize); - dsmPropagator.inject(span.with(dsmContext), record, SR_SETTER); - } - } + reportDsmCheckpointOrInject(span, record, tags, payloadSize); CONSUMER_DECORATE.afterStart(span); CONSUMER_DECORATE.onConsume(span, record, node); - AgentScope agentScope = activateSpan(span); if (null != queueSpan) { queueSpan.finish(); } + return span; + } - if (streamTaskContext == null) { - streamTaskContext = new StreamTaskContext(); + /** + * DSM-only mode (tracing disabled for kafka-streams, DSM enabled): never creates a real span, + * so no span is ever written to the agent for this integration. Only the pathway + * checkpoint/injection happens, carried by a lightweight, never-collected span shim. + */ + private static AgentSpan startDsmOnlyPathwaySpan( + final StampedRecord record, final String applicationId) { + final AgentSpan localActiveSpan = activeSpan(); + final AgentSpan span = + newPathwayOnlySpan(localActiveSpan == null ? null : localActiveSpan.spanContext()); + + DataStreamsTags tags = createWithGroup("kafka", INBOUND, applicationId, record.topic()); + final long payloadSize = computePayloadSizeBytes(record.value); + reportDsmCheckpointOrInject(span, record, tags, payloadSize); + return span; + } + + /** + * Reports a DSM checkpoint for {@code record}'s topic, or - when in a streaming context and + * {@code record}'s topic is a source topic - injects the pathway context so it survives leaving + * the topology on another instance of the application. + */ + private static void reportDsmCheckpointOrInject( + final AgentSpan span, + final StampedRecord record, + final DataStreamsTags tags, + final long payloadSize) { + if (STREAMING_CONTEXT.isDisabledForTopic(record.topic())) { + AgentTracer.get() + .getDataStreamsMonitoring() + .setCheckpoint(span, create(tags, record.timestamp, payloadSize)); + } else if (STREAMING_CONTEXT.isSourceTopic(record.topic())) { + Propagator dsmPropagator = Propagators.forConcern(DSM_CONCERN); + DataStreamsContext dsmContext = create(tags, record.timestamp, payloadSize); + dsmPropagator.inject(span.with(dsmContext), record, SR_SETTER); } - streamTaskContext.setAgentScope(agentScope); - InstrumentationContext.get(StreamTask.class, StreamTaskContext.class) - .put(task, streamTaskContext); } } @@ -342,9 +385,43 @@ public static void start( return; } - AgentSpan span, queueSpan = null; StreamTaskContext streamTaskContext = InstrumentationContext.get(StreamTask.class, StreamTaskContext.class).get(task); + String applicationId = + streamTaskContext != null ? streamTaskContext.getApplicationId() : null; + + final AgentSpan span = + !KafkaStreamsDecorator.TRACING_ENABLED && traceConfig().isDataStreamsEnabled() + ? startDsmOnlyPathwaySpan(record, applicationId) + : startTracedConsumeSpan(record, node, applicationId); + + AgentScope agentScope = activateSpan(span); + + if (streamTaskContext == null) { + streamTaskContext = new StreamTaskContext(); + } + streamTaskContext.setAgentScope(agentScope); + InstrumentationContext.get(StreamTask.class, StreamTaskContext.class) + .put(task, streamTaskContext); + } + + private static long payloadSizeBytes(final ProcessorRecordContext record) { + // we have to go through Object to get the RecordMetadata here because the class of `record` + // only implements it after 2.7 (and this class is only used if v >= 2.7) + if ((Object) record instanceof RecordMetadata) { // should always be true + RecordMetadata metadata = (RecordMetadata) (Object) record; + return metadata.serializedKeySize() + metadata.serializedValueSize(); + } + return 0; + } + + /** + * Creates and activates the real APM consume span (and, when time-in-queue is enabled, its + * broker parent), tags it, and reports DSM checkpoints/transactions off of it. + */ + private static AgentSpan startTracedConsumeSpan( + final ProcessorRecordContext record, final ProcessorNode node, final String applicationId) { + AgentSpan span, queueSpan = null; long timeInQueueStart = PR_GETTER.extractTimeInQueueStart(record); if (timeInQueueStart == 0 || !TIME_IN_QUEUE_ENABLED) { span = startSpan(JAVA_KAFKA.toString(), KAFKA_CONSUME); @@ -360,45 +437,54 @@ public static void start( // spans are written out together by TraceStructureWriter when running in strict mode } - String applicationId = null; - if (streamTaskContext != null) { - applicationId = streamTaskContext.getApplicationId(); - } DataStreamsTags tags = createWithGroup("kafka", INBOUND, applicationId, record.topic()); - - long payloadSize = 0; - // we have to go through Object to get the RecordMetadata here because the class of `record` - // only implements it after 2.7 (and this class is only used if v >= 2.7) - if ((Object) record instanceof RecordMetadata) { // should always be true - RecordMetadata metadata = (RecordMetadata) (Object) record; - payloadSize = metadata.serializedKeySize() + metadata.serializedValueSize(); - } - - if (STREAMING_CONTEXT.isDisabledForTopic(record.topic())) { - AgentTracer.get() - .getDataStreamsMonitoring() - .setCheckpoint(span, create(tags, record.timestamp(), payloadSize)); - } else { - if (STREAMING_CONTEXT.isSourceTopic(record.topic())) { - Propagator dsmPropagator = Propagators.forConcern(DSM_CONCERN); - DataStreamsContext dsmContext = create(tags, record.timestamp(), payloadSize); - dsmPropagator.inject(span.with(dsmContext), record, PR_SETTER); - } - } + long payloadSize = payloadSizeBytes(record); + reportDsmCheckpointOrInject(span, record, tags, payloadSize); CONSUMER_DECORATE.afterStart(span); CONSUMER_DECORATE.onConsume(span, record, node); - AgentScope agentScope = activateSpan(span); if (null != queueSpan) { queueSpan.finish(); } + return span; + } - if (streamTaskContext == null) { - streamTaskContext = new StreamTaskContext(); + /** + * DSM-only mode (tracing disabled for kafka-streams, DSM enabled): never creates a real span, + * so no span is ever written to the agent for this integration. Only the pathway + * checkpoint/injection happens, carried by a lightweight, never-collected span shim. + */ + private static AgentSpan startDsmOnlyPathwaySpan( + final ProcessorRecordContext record, final String applicationId) { + final AgentSpan localActiveSpan = activeSpan(); + final AgentSpan span = + newPathwayOnlySpan(localActiveSpan == null ? null : localActiveSpan.spanContext()); + + DataStreamsTags tags = createWithGroup("kafka", INBOUND, applicationId, record.topic()); + long payloadSize = payloadSizeBytes(record); + reportDsmCheckpointOrInject(span, record, tags, payloadSize); + return span; + } + + /** + * Reports a DSM checkpoint for {@code record}'s topic, or - when in a streaming context and + * {@code record}'s topic is a source topic - injects the pathway context so it survives leaving + * the topology on another instance of the application. + */ + private static void reportDsmCheckpointOrInject( + final AgentSpan span, + final ProcessorRecordContext record, + final DataStreamsTags tags, + final long payloadSize) { + if (STREAMING_CONTEXT.isDisabledForTopic(record.topic())) { + AgentTracer.get() + .getDataStreamsMonitoring() + .setCheckpoint(span, create(tags, record.timestamp(), payloadSize)); + } else if (STREAMING_CONTEXT.isSourceTopic(record.topic())) { + Propagator dsmPropagator = Propagators.forConcern(DSM_CONCERN); + DataStreamsContext dsmContext = create(tags, record.timestamp(), payloadSize); + dsmPropagator.inject(span.with(dsmContext), record, PR_SETTER); } - streamTaskContext.setAgentScope(agentScope); - InstrumentationContext.get(StreamTask.class, StreamTaskContext.class) - .put(task, streamTaskContext); } } diff --git a/dd-java-agent/instrumentation/kafka/kafka-streams-0.11/src/main/java/datadog/trace/instrumentation/kafka_streams/KafkaStreamsDecorator.java b/dd-java-agent/instrumentation/kafka/kafka-streams-0.11/src/main/java/datadog/trace/instrumentation/kafka_streams/KafkaStreamsDecorator.java index 97f52bae52c..99a505a0d63 100644 --- a/dd-java-agent/instrumentation/kafka/kafka-streams-0.11/src/main/java/datadog/trace/instrumentation/kafka_streams/KafkaStreamsDecorator.java +++ b/dd-java-agent/instrumentation/kafka/kafka-streams-0.11/src/main/java/datadog/trace/instrumentation/kafka_streams/KafkaStreamsDecorator.java @@ -7,6 +7,7 @@ import datadog.trace.api.Config; import datadog.trace.api.Functions; +import datadog.trace.api.InstrumenterConfig; import datadog.trace.api.cache.DDCache; import datadog.trace.api.cache.DDCaches; import datadog.trace.api.naming.SpanNaming; @@ -15,6 +16,7 @@ import datadog.trace.bootstrap.instrumentation.api.Tags; import datadog.trace.bootstrap.instrumentation.api.UTF8BytesString; import datadog.trace.bootstrap.instrumentation.decorator.MessagingClientDecorator; +import java.util.Arrays; import java.util.function.Supplier; import org.apache.kafka.streams.processor.internals.ProcessorNode; import org.apache.kafka.streams.processor.internals.ProcessorRecordContext; @@ -22,6 +24,10 @@ public class KafkaStreamsDecorator extends MessagingClientDecorator { private static final String KAFKA = "kafka"; + // Kept in sync with the names each kafka-streams-0.11 instrumentation module passes to its own + // super(...) constructor call, so TRACING_ENABLED can't drift from what is actually registered. + public static final String INTEGRATION_NAME = KAFKA; + public static final String LEGACY_INTEGRATION_NAME = "kafka-streams"; public static final CharSequence JAVA_KAFKA = UTF8BytesString.create("java-kafka-streams"); public static final CharSequence KAFKA_CONSUME = UTF8BytesString.create( @@ -31,6 +37,11 @@ public class KafkaStreamsDecorator extends MessagingClientDecorator { public static final boolean KAFKA_LEGACY_TRACING = Config.get().isKafkaLegacyTracingEnabled(); public static final boolean TIME_IN_QUEUE_ENABLED = Config.get().isTimeInQueueEnabled(!KAFKA_LEGACY_TRACING, KAFKA); + public static final boolean TRACING_ENABLED = + InstrumenterConfig.get() + .isIntegrationEnabled( + Arrays.asList(INTEGRATION_NAME, LEGACY_INTEGRATION_NAME), + InstrumenterConfig.get().isIntegrationsEnabled()); public static final String KAFKA_PRODUCED_KEY = "x_datadog_kafka_produced"; private final String spanKind; diff --git a/dd-java-agent/instrumentation/kafka/kafka-streams-0.11/src/main/java/datadog/trace/instrumentation/kafka_streams/KafkaStreamsSourceNodeRecordDeserializerInstrumentation.java b/dd-java-agent/instrumentation/kafka/kafka-streams-0.11/src/main/java/datadog/trace/instrumentation/kafka_streams/KafkaStreamsSourceNodeRecordDeserializerInstrumentation.java index 81348490676..f2c473591d1 100644 --- a/dd-java-agent/instrumentation/kafka/kafka-streams-0.11/src/main/java/datadog/trace/instrumentation/kafka_streams/KafkaStreamsSourceNodeRecordDeserializerInstrumentation.java +++ b/dd-java-agent/instrumentation/kafka/kafka-streams-0.11/src/main/java/datadog/trace/instrumentation/kafka_streams/KafkaStreamsSourceNodeRecordDeserializerInstrumentation.java @@ -16,11 +16,11 @@ // This is necessary because SourceNodeRecordDeserializer drops the headers. :-( @AutoService(InstrumenterModule.class) public class KafkaStreamsSourceNodeRecordDeserializerInstrumentation - extends InstrumenterModule.Tracing + extends InstrumenterModule.DataStreams implements Instrumenter.ForSingleType, Instrumenter.HasMethodAdvice { public KafkaStreamsSourceNodeRecordDeserializerInstrumentation() { - super("kafka", "kafka-streams"); + super(KafkaStreamsDecorator.INTEGRATION_NAME, KafkaStreamsDecorator.LEGACY_INTEGRATION_NAME); } @Override diff --git a/dd-java-agent/instrumentation/kafka/kafka-streams-0.11/src/test/groovy/KafkaStreamsDataStreamsOnlyForkedTest.groovy b/dd-java-agent/instrumentation/kafka/kafka-streams-0.11/src/test/groovy/KafkaStreamsDataStreamsOnlyForkedTest.groovy new file mode 100644 index 00000000000..ad3d2ff0f27 --- /dev/null +++ b/dd-java-agent/instrumentation/kafka/kafka-streams-0.11/src/test/groovy/KafkaStreamsDataStreamsOnlyForkedTest.groovy @@ -0,0 +1,151 @@ +import datadog.trace.agent.test.InstrumentationSpecification +import datadog.trace.api.config.TraceInstrumentationConfig +import org.apache.kafka.clients.producer.KafkaProducer +import org.apache.kafka.clients.producer.ProducerRecord +import org.apache.kafka.common.header.internals.RecordHeader +import org.apache.kafka.common.header.internals.RecordHeaders +import org.apache.kafka.common.serialization.Serdes +import org.apache.kafka.common.serialization.StringSerializer +import org.apache.kafka.streams.KafkaStreams +import org.apache.kafka.streams.StreamsConfig +import org.apache.kafka.streams.kstream.KStream +import org.apache.kafka.streams.kstream.KStreamBuilder +import org.apache.kafka.streams.kstream.ValueMapper +import org.springframework.kafka.test.rule.KafkaEmbedded +import org.springframework.kafka.test.utils.KafkaTestUtils +import spock.lang.Shared + +import java.nio.charset.StandardCharsets + +/** + * DSM-only coverage for the kafka-streams StreamTask consume path: when Kafka APM tracing is + * disabled (integrations.enabled=false) but DSM is enabled, this integration must never create or + * write a real span/trace, whether or not the record carries a propagated trace context. + */ +abstract class KafkaStreamsDataStreamsOnlyForkedTest extends InstrumentationSpecification { + static final STREAM_PENDING = "test.pending" + static final STREAM_PROCESSED = "test.processed" + + @Shared + protected KafkaEmbedded embeddedKafka + + def setupSpec() { + embeddedKafka = new KafkaEmbedded(1, true, 1, STREAM_PENDING, STREAM_PROCESSED) + embeddedKafka.before() + } + + def cleanupSpec() { + embeddedKafka?.after() + } + + @Override + void configurePreAgent() { + super.configurePreAgent() + injectSysConfig("integrations.enabled", "false") + } + + @Override + boolean useStrictTraceWrites() { + return false + } + + @Override + protected boolean isDataStreamsEnabled() { + return true + } + + protected KafkaStreams startLowercasingTopology() { + def config = new Properties() + config.putAll(KafkaTestUtils.senderProps(embeddedKafka.getBrokersAsString())) + config.put(StreamsConfig.APPLICATION_ID_CONFIG, "dsm-only-test-application") + config.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass().getName()) + config.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.String().getClass().getName()) + + def builder = new KStreamBuilder() + KStream textLines = builder.stream(STREAM_PENDING) + textLines + .mapValues(new ValueMapper() { + @Override + String apply(String textLine) { + return textLine.toLowerCase() + } + }) + .to(Serdes.String(), Serdes.String(), STREAM_PROCESSED) + + def streams = new KafkaStreams(builder, config) + streams.start() + return streams + } + + protected KafkaProducer newProducer() { + return new KafkaProducer( + KafkaTestUtils.senderProps(embeddedKafka.getBrokersAsString()), + new StringSerializer(), + new StringSerializer()) + } +} + +/** + * A record with no propagated Datadog trace context must not create any streams consume span in + * DSM-only mode - only a DSM checkpoint, tracked via the lightweight pathway-only span shim. + */ +class KafkaStreamsDataStreamsOnlyLocalRootForkedTest extends KafkaStreamsDataStreamsOnlyForkedTest { + + @Override + void configurePreAgent() { + super.configurePreAgent() + // The in-JVM producer would otherwise inject its own trace context into the record headers, + // which the streams ContextPropagationAdvice would then extract - so there would be no way to + // exercise the genuinely-local-root path. Disabling client propagation for the topic stops + // both the injection and the extraction. + injectSysConfig(TraceInstrumentationConfig.KAFKA_CLIENT_PROPAGATION_DISABLED_TOPICS, STREAM_PENDING) + } + + def "a local-root record does not create a streams consume span"() { + setup: + def streams = startLowercasingTopology() + def producer = newProducer() + + when: + producer.send(new ProducerRecord(STREAM_PENDING, "LOCAL ROOT")).get() + + then: + TEST_DATA_STREAMS_WRITER.waitForGroups(1) + TEST_WRITER.isEmpty() + + cleanup: + producer?.close() + streams?.close() + } +} + +/** + * Regression guard: a record carrying a real, externally-propagated Datadog trace context must + * still not create any streams consume span in DSM-only mode. + */ +class KafkaStreamsDataStreamsOnlyExtractedParentForkedTest extends KafkaStreamsDataStreamsOnlyForkedTest { + + def "a record continuing a propagated trace does not create a streams consume span either"() { + setup: + def streams = startLowercasingTopology() + def producer = newProducer() + def existingTraceId = 1234567890123456L + def existingSpanId = 9876543210987654L + def headers = new RecordHeaders() + headers.add(new RecordHeader("x-datadog-trace-id", + String.valueOf(existingTraceId).getBytes(StandardCharsets.UTF_8))) + headers.add(new RecordHeader("x-datadog-parent-id", + String.valueOf(existingSpanId).getBytes(StandardCharsets.UTF_8))) + + when: + producer.send(new ProducerRecord(STREAM_PENDING, null, null, "PROPAGATED", headers)).get() + + then: + TEST_DATA_STREAMS_WRITER.waitForGroups(1) + TEST_WRITER.isEmpty() + + cleanup: + producer?.close() + streams?.close() + } +} diff --git a/dd-java-agent/instrumentation/kafka/kafka-streams-1.0/src/main/java/datadog/trace/instrumentation/kafka_streams10/InternalTopologyBuilderInstrumentation.java b/dd-java-agent/instrumentation/kafka/kafka-streams-1.0/src/main/java/datadog/trace/instrumentation/kafka_streams10/InternalTopologyBuilderInstrumentation.java index 07271bb7904..a32e48ac094 100644 --- a/dd-java-agent/instrumentation/kafka/kafka-streams-1.0/src/main/java/datadog/trace/instrumentation/kafka_streams10/InternalTopologyBuilderInstrumentation.java +++ b/dd-java-agent/instrumentation/kafka/kafka-streams-1.0/src/main/java/datadog/trace/instrumentation/kafka_streams10/InternalTopologyBuilderInstrumentation.java @@ -12,7 +12,7 @@ import org.apache.kafka.streams.processor.internals.ProcessorTopology; @AutoService(InstrumenterModule.class) -public class InternalTopologyBuilderInstrumentation extends InstrumenterModule.Tracing +public class InternalTopologyBuilderInstrumentation extends InstrumenterModule.DataStreams implements Instrumenter.ForSingleType, Instrumenter.HasMethodAdvice { public InternalTopologyBuilderInstrumentation() { diff --git a/internal-api/src/main/java/datadog/trace/api/Config.java b/internal-api/src/main/java/datadog/trace/api/Config.java index fade2b4c417..7af8f65207d 100644 --- a/internal-api/src/main/java/datadog/trace/api/Config.java +++ b/internal-api/src/main/java/datadog/trace/api/Config.java @@ -53,7 +53,6 @@ import static datadog.trace.api.ConfigDefaults.DEFAULT_DATA_JOBS_OPENLINEAGE_TIMEOUT_ENABLED; import static datadog.trace.api.ConfigDefaults.DEFAULT_DATA_JOBS_PARSE_SPARK_PLAN_ENABLED; import static datadog.trace.api.ConfigDefaults.DEFAULT_DATA_STREAMS_BUCKET_DURATION; -import static datadog.trace.api.ConfigDefaults.DEFAULT_DATA_STREAMS_ENABLED; import static datadog.trace.api.ConfigDefaults.DEFAULT_DB_CLIENT_HOST_SPLIT_BY_HOST; import static datadog.trace.api.ConfigDefaults.DEFAULT_DB_CLIENT_HOST_SPLIT_BY_INSTANCE; import static datadog.trace.api.ConfigDefaults.DEFAULT_DB_CLIENT_HOST_SPLIT_BY_INSTANCE_TYPE_SUFFIX; @@ -382,7 +381,6 @@ import static datadog.trace.api.config.GeneralConfig.DATA_JOBS_OPENLINEAGE_TIMEOUT_ENABLED; import static datadog.trace.api.config.GeneralConfig.DATA_JOBS_PARSE_SPARK_PLAN_ENABLED; import static datadog.trace.api.config.GeneralConfig.DATA_STREAMS_BUCKET_DURATION_SECONDS; -import static datadog.trace.api.config.GeneralConfig.DATA_STREAMS_ENABLED; import static datadog.trace.api.config.GeneralConfig.DATA_STREAMS_TRANSACTION_EXTRACTORS; import static datadog.trace.api.config.GeneralConfig.DOGSTATSD_ARGS; import static datadog.trace.api.config.GeneralConfig.DOGSTATSD_HOST; @@ -1367,7 +1365,6 @@ public static String getHostName() { private final boolean dataJobsParseSparkPlanEnabled; private final boolean dataJobsExperimentalFeaturesEnabled; - private final boolean dataStreamsEnabled; private final float dataStreamsBucketDurationSeconds; private final String dataStreamsTransactionExtractors; @@ -3192,8 +3189,6 @@ PROFILING_DATADOG_PROFILER_ENABLED, isDatadogProfilerSafeInCurrentEnvironment()) DATA_JOBS_EXPERIMENTAL_FEATURES_ENABLED, DEFAULT_DATA_JOBS_EXPERIMENTAL_FEATURES_ENABLED); - dataStreamsEnabled = - configProvider.getBoolean(DATA_STREAMS_ENABLED, DEFAULT_DATA_STREAMS_ENABLED); dataStreamsBucketDurationSeconds = configProvider.getFloat( DATA_STREAMS_BUCKET_DURATION_SECONDS, DEFAULT_DATA_STREAMS_BUCKET_DURATION); @@ -5169,7 +5164,7 @@ public boolean isAwsServerless() { } public boolean isDataStreamsEnabled() { - return dataStreamsEnabled; + return instrumenterConfig.isDataStreamsEnabled(); } public float getDataStreamsBucketDurationSeconds() { diff --git a/internal-api/src/main/java/datadog/trace/api/InstrumenterConfig.java b/internal-api/src/main/java/datadog/trace/api/InstrumenterConfig.java index 74bff640024..dd32dff1c6f 100644 --- a/internal-api/src/main/java/datadog/trace/api/InstrumenterConfig.java +++ b/internal-api/src/main/java/datadog/trace/api/InstrumenterConfig.java @@ -7,6 +7,7 @@ import static datadog.trace.api.ConfigDefaults.DEFAULT_CIVISIBILITY_ENABLED; import static datadog.trace.api.ConfigDefaults.DEFAULT_CODE_ORIGIN_FOR_SPANS_INTERFACE_SUPPORT; import static datadog.trace.api.ConfigDefaults.DEFAULT_DATA_JOBS_ENABLED; +import static datadog.trace.api.ConfigDefaults.DEFAULT_DATA_STREAMS_ENABLED; import static datadog.trace.api.ConfigDefaults.DEFAULT_IAST_ENABLED; import static datadog.trace.api.ConfigDefaults.DEFAULT_INTEGRATIONS_ENABLED; import static datadog.trace.api.ConfigDefaults.DEFAULT_LLM_OBS_ENABLED; @@ -35,6 +36,7 @@ import static datadog.trace.api.config.GeneralConfig.AGENTLESS_LOG_SUBMISSION_ENABLED; import static datadog.trace.api.config.GeneralConfig.APP_LOGS_COLLECTION_ENABLED; import static datadog.trace.api.config.GeneralConfig.DATA_JOBS_ENABLED; +import static datadog.trace.api.config.GeneralConfig.DATA_STREAMS_ENABLED; import static datadog.trace.api.config.GeneralConfig.INTERNAL_EXIT_ON_FAILURE; import static datadog.trace.api.config.GeneralConfig.TELEMETRY_ENABLED; import static datadog.trace.api.config.GeneralConfig.TRACE_DEBUG; @@ -161,6 +163,7 @@ public class InstrumenterConfig { private final boolean appSecRaspEnabled; private final boolean iastFullyDisabled; private final boolean usmEnabled; + private final boolean dataStreamsEnabled; private final boolean telemetryEnabled; private final boolean llmObsEnabled; @@ -287,6 +290,8 @@ private InstrumenterConfig() { final Boolean iastEnabled = configProvider.getBoolean(IAST_ENABLED); iastFullyDisabled = iastEnabled != null && !iastEnabled; usmEnabled = configProvider.getBoolean(USM_ENABLED, DEFAULT_USM_ENABLED); + dataStreamsEnabled = + configProvider.getBoolean(DATA_STREAMS_ENABLED, DEFAULT_DATA_STREAMS_ENABLED); telemetryEnabled = configProvider.getBoolean(TELEMETRY_ENABLED, DEFAULT_TELEMETRY_ENABLED); llmObsEnabled = configProvider.getBoolean(LLMOBS_ENABLED, DEFAULT_LLM_OBS_ENABLED); } else { @@ -297,6 +302,7 @@ private InstrumenterConfig() { iastFullyDisabled = true; telemetryEnabled = false; usmEnabled = false; + dataStreamsEnabled = false; llmObsEnabled = false; } @@ -505,6 +511,10 @@ public boolean isUsmEnabled() { return usmEnabled; } + public boolean isDataStreamsEnabled() { + return dataStreamsEnabled; + } + public boolean isTelemetryEnabled() { return telemetryEnabled; } diff --git a/internal-api/src/main/java/datadog/trace/api/sampling/SamplingMechanism.java b/internal-api/src/main/java/datadog/trace/api/sampling/SamplingMechanism.java index 6b0cd397626..25fb473cca1 100644 --- a/internal-api/src/main/java/datadog/trace/api/sampling/SamplingMechanism.java +++ b/internal-api/src/main/java/datadog/trace/api/sampling/SamplingMechanism.java @@ -41,6 +41,7 @@ public class SamplingMechanism { public static final byte REMOTE_USER_RULE = 11; public static final byte REMOTE_ADAPTIVE_RULE = 12; public static final byte AI_GUARD = 13; + public static final byte DATA_STREAMS = 14; /** Force override sampling decision from external source, like W3C traceparent. */ public static final byte EXTERNAL_OVERRIDE = Byte.MIN_VALUE; @@ -68,6 +69,9 @@ public static boolean validateWithSamplingPriority(int mechanism, int priority) case DATA_JOBS: return priority == PrioritySampling.USER_KEEP; + case DATA_STREAMS: + return priority == PrioritySampling.USER_DROP; + case EXTERNAL_OVERRIDE: return false; } @@ -83,7 +87,8 @@ public static boolean validateWithSamplingPriority(int mechanism, int priority) */ public static boolean canAvoidSamplingPriorityLock(int priority, int mechanism) { return (!Config.get().isApmTracingEnabled() && mechanism == SamplingMechanism.APPSEC) - || (Config.get().isDataJobsEnabled() && mechanism == DATA_JOBS); + || (Config.get().isDataJobsEnabled() && mechanism == DATA_JOBS) + || (Config.get().isDataStreamsEnabled() && mechanism == DATA_STREAMS); } private SamplingMechanism() {} diff --git a/internal-api/src/test/groovy/datadog/trace/api/InstrumenterConfigTest.groovy b/internal-api/src/test/groovy/datadog/trace/api/InstrumenterConfigTest.groovy deleted file mode 100644 index 15c15c1ef3d..00000000000 --- a/internal-api/src/test/groovy/datadog/trace/api/InstrumenterConfigTest.groovy +++ /dev/null @@ -1,195 +0,0 @@ -package datadog.trace.api - -import datadog.trace.config.inversion.ConfigHelper -import datadog.trace.test.util.DDSpecification - -class InstrumenterConfigTest extends DDSpecification { - - def strictness - - def setup(){ - strictness = ConfigHelper.get().configInversionStrictFlag() - ConfigHelper.get().setConfigInversionStrict(ConfigHelper.StrictnessPolicy.TEST) - } - - def cleanup(){ - ConfigHelper.get().setConfigInversionStrict(strictness) - } - - def "verify integration config"() { - setup: - environmentVariables.set("DD_INTEGRATION_ORDER_ENABLED", "false") - environmentVariables.set("DD_INTEGRATION_TEST_ENV_ENABLED", "true") - environmentVariables.set("DD_TRACE_NEW_ENV_ENABLED", "false") - environmentVariables.set("DD_INTEGRATION_DISABLED_ENV_ENABLED", "false") - - System.setProperty("dd.integration.order.enabled", "true") - System.setProperty("dd.integration.test-prop.enabled", "true") - System.setProperty("dd.integration.disabled-prop.enabled", "false") - - environmentVariables.set("DD_INTEGRATION_ORDER_MATCHING_SHORTCUT_ENABLED", "false") - environmentVariables.set("DD_INTEGRATION_TEST_ENV_MATCHING_SHORTCUT_ENABLED", "true") - environmentVariables.set("DD_INTEGRATION_NEW_ENV_MATCHING_SHORTCUT_ENABLED", "false") - environmentVariables.set("DD_INTEGRATION_DISABLED_ENV_MATCHING_SHORTCUT_ENABLED", "false") - - System.setProperty("dd.integration.order.matching.shortcut.enabled", "true") - System.setProperty("dd.integration.test-prop.matching.shortcut.enabled", "true") - System.setProperty("dd.integration.disabled-prop.matching.shortcut.enabled", "false") - - expect: - InstrumenterConfig.get().isIntegrationEnabled(integrationNames, defaultEnabled) == expected - InstrumenterConfig.get().isIntegrationShortcutMatchingEnabled(integrationNames, defaultEnabled) == expected - - where: - // spotless:off - names | defaultEnabled | expected - [] | true | true - [] | false | false - ["invalid"] | true | true - ["invalid"] | false | false - ["test-prop"] | false | true - ["test-env"] | false | true - ["disabled-prop"] | true | false - ["disabled-env"] | true | false - ["other", "test-prop"] | false | true - ["other", "test-env"] | false | true - ["order"] | false | true - ["test-prop", "disabled-prop"] | false | true - ["disabled-env", "test-env"] | false | true - ["test-prop", "disabled-prop"] | true | false - ["disabled-env", "test-env"] | true | false - ["new-env"] | true | false - // spotless:on - - integrationNames = new TreeSet<>(names) - } - - def setEnv(String key, String value) { - environmentVariables.set(key, value) - } - - def setSysProp(String key, String value) { - System.setProperty(key, value) - } - - def randomIntegrationEnabled() { - return InstrumenterConfig.get().isIntegrationEnabled(["random"], true) - } - - def "verify integration enabled hierarchy"() { - when: - // the below should have no effect - setEnv("DD_RANDOM_ENABLED", "false") - setSysProp("dd.random.enabled", "false") - - then: - randomIntegrationEnabled() == true - - when: - setEnv("DD_INTEGRATION_RANDOM_ENABLED", "false") - - then: - randomIntegrationEnabled() == false - - when: - setEnv("DD_TRACE_INTEGRATION_RANDOM_ENABLED", "true") - - then: - randomIntegrationEnabled() == true - - when: - setEnv("DD_TRACE_RANDOM_ENABLED", "false") - - then: - randomIntegrationEnabled() == false - - // assert all system properties take precedence over all env vars - when: - setSysProp("dd.integration.random.enabled", "true") - - then: - randomIntegrationEnabled() == true - - when: - setSysProp("dd.trace.integration.random.enabled", "false") - - then: - randomIntegrationEnabled() == false - - when: - setSysProp("dd.trace.random.enabled", "true") - - then: - randomIntegrationEnabled() == true - } - - def "valid resolver presets"() { - setup: - injectSysConfig("resolver.cache.config", preset) - - expect: - InstrumenterConfig.get().resolverOutliningEnabled == outlining - - where: - // spotless:off - preset | outlining - 'LARGE' | true - 'SMALL' | true - 'DEFAULT' | true - 'LEGACY' | false - // spotless:on - } - - def "invalid resolver presets"() { - setup: - injectSysConfig("resolver.cache.config", preset) - - expect: - InstrumenterConfig.get().resolverOutliningEnabled - - where: - preset << ['INVALID', ''] - } - - def "appsec enabled = #input"() { - setup: - if (input != null) { - injectSysConfig("appsec.enabled", input) - } - - expect: - InstrumenterConfig.get().getAppSecActivation() == expected - - where: - input | expected - null | ProductActivation.ENABLED_INACTIVE - "" | ProductActivation.ENABLED_INACTIVE - "bad" | ProductActivation.FULLY_DISABLED - "false" | ProductActivation.FULLY_DISABLED - "0" | ProductActivation.FULLY_DISABLED - "true" | ProductActivation.FULLY_ENABLED - "1" | ProductActivation.FULLY_ENABLED - "inactive" | ProductActivation.ENABLED_INACTIVE - } - - def "iast enabled = #input"() { - setup: - if (input != null) { - injectSysConfig("iast.enabled", input) - } - - expect: - InstrumenterConfig.get().getIastActivation() == expected - - where: - input | expected - null | ProductActivation.FULLY_DISABLED - "" | ProductActivation.FULLY_DISABLED - "bad" | ProductActivation.FULLY_DISABLED - "false" | ProductActivation.FULLY_DISABLED - "0" | ProductActivation.FULLY_DISABLED - "true" | ProductActivation.FULLY_ENABLED - "1" | ProductActivation.FULLY_ENABLED - "inactive" | ProductActivation.ENABLED_INACTIVE - } -} diff --git a/internal-api/src/test/groovy/datadog/trace/api/sampling/SamplingMechanismTest.groovy b/internal-api/src/test/groovy/datadog/trace/api/sampling/SamplingMechanismTest.groovy deleted file mode 100644 index 4a4890435c1..00000000000 --- a/internal-api/src/test/groovy/datadog/trace/api/sampling/SamplingMechanismTest.groovy +++ /dev/null @@ -1,120 +0,0 @@ -package datadog.trace.api.sampling - -import datadog.trace.test.util.DDSpecification -import static datadog.trace.api.sampling.PrioritySampling.* -import static datadog.trace.api.sampling.SamplingMechanism.* - -class SamplingMechanismTest extends DDSpecification { - - static userDropX = USER_DROP - 1 - static userKeepX = USER_KEEP + 1 - - def "test validation"() { - expect: - validateWithSamplingPriority(mechanism, priority) == valid - - where: - mechanism | priority | valid - UNKNOWN | UNSET | true - UNKNOWN | SAMPLER_DROP | true - UNKNOWN | SAMPLER_KEEP | true - UNKNOWN | USER_DROP | true - UNKNOWN | USER_KEEP | true - UNKNOWN | userDropX | true - UNKNOWN | userKeepX | true - - DEFAULT | UNSET | false - DEFAULT | SAMPLER_DROP | true - DEFAULT | SAMPLER_KEEP | true - DEFAULT | USER_DROP | false - DEFAULT | USER_KEEP | false - DEFAULT | userDropX | false - DEFAULT | userKeepX | false - - AGENT_RATE | UNSET | false - AGENT_RATE | SAMPLER_DROP | true - AGENT_RATE | SAMPLER_KEEP | true - AGENT_RATE | USER_DROP | false - AGENT_RATE | USER_KEEP | false - AGENT_RATE | userDropX | false - AGENT_RATE | userKeepX | false - - REMOTE_AUTO_RATE | UNSET | false - REMOTE_AUTO_RATE | SAMPLER_DROP | true - REMOTE_AUTO_RATE | SAMPLER_KEEP | true - REMOTE_AUTO_RATE | USER_DROP | false - REMOTE_AUTO_RATE | USER_KEEP | false - REMOTE_AUTO_RATE | userDropX | false - REMOTE_AUTO_RATE | userKeepX | false - - LOCAL_USER_RULE | UNSET | false - LOCAL_USER_RULE | SAMPLER_DROP | false - LOCAL_USER_RULE | SAMPLER_KEEP | false - LOCAL_USER_RULE | USER_DROP | true - LOCAL_USER_RULE | USER_KEEP | true - LOCAL_USER_RULE | userDropX | false - LOCAL_USER_RULE | userKeepX | false - - MANUAL | UNSET | false - MANUAL | SAMPLER_DROP | false - MANUAL | SAMPLER_KEEP | false - MANUAL | USER_DROP | true - MANUAL | USER_KEEP | true - MANUAL | userDropX | false - MANUAL | userKeepX | false - - REMOTE_USER_RATE | UNSET | false - REMOTE_USER_RATE | SAMPLER_DROP | false - REMOTE_USER_RATE | SAMPLER_KEEP | false - REMOTE_USER_RATE | USER_DROP | true - REMOTE_USER_RATE | USER_KEEP | true - REMOTE_USER_RATE | userDropX | false - REMOTE_USER_RATE | userKeepX | false - - APPSEC | UNSET | false - APPSEC | SAMPLER_DROP | true - APPSEC | SAMPLER_KEEP | true - APPSEC | USER_DROP | false - APPSEC | USER_KEEP | true - APPSEC | userDropX | false - APPSEC | userKeepX | false - - DATA_JOBS | UNSET | false - DATA_JOBS | SAMPLER_DROP | false - DATA_JOBS | SAMPLER_KEEP | false - DATA_JOBS | USER_DROP | false - DATA_JOBS | USER_KEEP | true - DATA_JOBS | userDropX | false - DATA_JOBS | userKeepX | false - - EXTERNAL_OVERRIDE | UNSET | false - EXTERNAL_OVERRIDE | SAMPLER_DROP | false - EXTERNAL_OVERRIDE | SAMPLER_KEEP | false - EXTERNAL_OVERRIDE | USER_DROP | false - EXTERNAL_OVERRIDE | USER_KEEP | false - EXTERNAL_OVERRIDE | userDropX | false - EXTERNAL_OVERRIDE | userKeepX | false - } - - void 'Test canAvoidSamplingPriorityLock'(){ - setup: - injectSysConfig("dd.apm.tracing.enabled", "false") - - expect: - canAvoidSamplingPriorityLock(priority, mechanism) == valid - - where: - mechanism | priority | valid - APPSEC | UNSET | true - APPSEC | SAMPLER_KEEP | true - UNKNOWN | SAMPLER_KEEP | false - DEFAULT | SAMPLER_KEEP | false - AGENT_RATE | SAMPLER_KEEP | false - REMOTE_AUTO_RATE | SAMPLER_KEEP | false - LOCAL_USER_RULE | SAMPLER_KEEP | false - MANUAL | SAMPLER_KEEP | false - REMOTE_USER_RATE | SAMPLER_KEEP | false - DATA_JOBS | SAMPLER_KEEP | false - EXTERNAL_OVERRIDE | SAMPLER_KEEP | false - } -} diff --git a/internal-api/src/test/java/datadog/trace/api/InstrumenterConfigTest.java b/internal-api/src/test/java/datadog/trace/api/InstrumenterConfigTest.java new file mode 100644 index 00000000000..0916b2ca9f3 --- /dev/null +++ b/internal-api/src/test/java/datadog/trace/api/InstrumenterConfigTest.java @@ -0,0 +1,198 @@ +package datadog.trace.api; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import datadog.trace.config.inversion.ConfigHelper; +import datadog.trace.config.inversion.ConfigHelper.StrictnessPolicy; +import datadog.trace.test.junit.utils.config.WithConfig; +import datadog.trace.test.junit.utils.config.WithConfigExtension; +import java.util.Collections; +import java.util.List; +import java.util.Set; +import java.util.TreeSet; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +import org.tabletest.junit.TableTest; + +@ExtendWith(WithConfigExtension.class) +class InstrumenterConfigTest { + + private StrictnessPolicy strictness; + + @BeforeEach + void setup() { + strictness = ConfigHelper.get().configInversionStrictFlag(); + ConfigHelper.get().setConfigInversionStrict(StrictnessPolicy.TEST); + } + + @AfterEach + void cleanup() { + ConfigHelper.get().setConfigInversionStrict(strictness); + } + + @TableTest({ + "scenario | names | defaultEnabled | expected", + "empty names, default enabled | [] | true | true ", + "empty names, default disabled | [] | false | false ", + "invalid name, default enabled | [invalid] | true | true ", + "invalid name, default disabled | [invalid] | false | false ", + "test-prop env var enabled overrides default off | [test-prop] | false | true ", + "test-env env var enabled overrides default off | [test-env] | false | true ", + "disabled-prop sys prop overrides default on | [disabled-prop] | true | false ", + "disabled-env env var overrides default on | [disabled-env] | true | false ", + "mixed names, test-prop wins | [other, test-prop] | false | true ", + "mixed names, test-env wins | [other, test-env] | false | true ", + "order enabled by both sys prop and env var | [order] | false | true ", + "test-prop and disabled-prop, default off | [test-prop, disabled-prop] | false | true ", + "disabled-env and test-env, default off | [disabled-env, test-env] | false | true ", + "test-prop and disabled-prop, default on | [test-prop, disabled-prop] | true | false ", + "disabled-env and test-env, default on | [disabled-env, test-env] | true | false ", + "new-env disabled overrides default on | [new-env] | true | false " + }) + @WithConfig(key = "INTEGRATION_ORDER_ENABLED", value = "false", env = true) + @WithConfig(key = "INTEGRATION_TEST_ENV_ENABLED", value = "true", env = true) + @WithConfig(key = "TRACE_NEW_ENV_ENABLED", value = "false", env = true) + @WithConfig(key = "INTEGRATION_DISABLED_ENV_ENABLED", value = "false", env = true) + @WithConfig(key = "INTEGRATION_ORDER_MATCHING_SHORTCUT_ENABLED", value = "false", env = true) + @WithConfig(key = "INTEGRATION_TEST_ENV_MATCHING_SHORTCUT_ENABLED", value = "true", env = true) + @WithConfig(key = "INTEGRATION_NEW_ENV_MATCHING_SHORTCUT_ENABLED", value = "false", env = true) + @WithConfig( + key = "INTEGRATION_DISABLED_ENV_MATCHING_SHORTCUT_ENABLED", + value = "false", + env = true) + @WithConfig(key = "integration.order.enabled", value = "true") + @WithConfig(key = "integration.test-prop.enabled", value = "true") + @WithConfig(key = "integration.disabled-prop.enabled", value = "false") + @WithConfig(key = "integration.order.matching.shortcut.enabled", value = "true") + @WithConfig(key = "integration.test-prop.matching.shortcut.enabled", value = "true") + @WithConfig(key = "integration.disabled-prop.matching.shortcut.enabled", value = "false") + void verifyIntegrationConfig(List names, boolean defaultEnabled, boolean expected) { + Set integrationNames = new TreeSet<>(names); + assertEquals( + expected, InstrumenterConfig.get().isIntegrationEnabled(integrationNames, defaultEnabled)); + assertEquals( + expected, + InstrumenterConfig.get() + .isIntegrationShortcutMatchingEnabled(integrationNames, defaultEnabled)); + } + + private static boolean randomIntegrationEnabled() { + return InstrumenterConfig.get().isIntegrationEnabled(Collections.singletonList("random"), true); + } + + @Test + void verifyIntegrationEnabledHierarchy() { + // the below should have no effect + WithConfigExtension.injectEnvConfig("RANDOM_ENABLED", "false"); + WithConfigExtension.injectSysConfig("random.enabled", "false"); + assertTrue(randomIntegrationEnabled()); + + WithConfigExtension.injectEnvConfig("INTEGRATION_RANDOM_ENABLED", "false"); + assertFalse(randomIntegrationEnabled()); + + WithConfigExtension.injectEnvConfig("TRACE_INTEGRATION_RANDOM_ENABLED", "true"); + assertTrue(randomIntegrationEnabled()); + + WithConfigExtension.injectEnvConfig("TRACE_RANDOM_ENABLED", "false"); + assertFalse(randomIntegrationEnabled()); + + // assert all system properties take precedence over all env vars + WithConfigExtension.injectSysConfig("integration.random.enabled", "true"); + assertTrue(randomIntegrationEnabled()); + + WithConfigExtension.injectSysConfig("trace.integration.random.enabled", "false"); + assertFalse(randomIntegrationEnabled()); + + WithConfigExtension.injectSysConfig("trace.random.enabled", "true"); + assertTrue(randomIntegrationEnabled()); + } + + @TableTest({ + "scenario | preset | outlining", + "large preset | LARGE | true ", + "small preset | SMALL | true ", + "default preset | DEFAULT | true ", + "legacy preset | LEGACY | false " + }) + void validResolverPresets(String preset, boolean outlining) { + WithConfigExtension.injectSysConfig("resolver.cache.config", preset); + + assertEquals(outlining, InstrumenterConfig.get().isResolverOutliningEnabled()); + } + + @TableTest({"scenario | preset ", "invalid preset | INVALID", "empty preset | '' "}) + void invalidResolverPresets(String preset) { + WithConfigExtension.injectSysConfig("resolver.cache.config", preset); + + assertTrue(InstrumenterConfig.get().isResolverOutliningEnabled()); + } + + @TableTest({ + "scenario | input | expected ", + "unset defaults to inactive | | ENABLED_INACTIVE", + "empty string is inactive | '' | ENABLED_INACTIVE", + "unparseable value disables | bad | FULLY_DISABLED ", + "explicit false disables | false | FULLY_DISABLED ", + "zero disables | 0 | FULLY_DISABLED ", + "explicit true enables | true | FULLY_ENABLED ", + "one enables | 1 | FULLY_ENABLED ", + "inactive keyword enables inactive | inactive | ENABLED_INACTIVE" + }) + void appsecEnabled(String input, ProductActivation expected) { + if (input != null) { + WithConfigExtension.injectSysConfig("appsec.enabled", input); + } + + assertEquals(expected, InstrumenterConfig.get().getAppSecActivation()); + } + + @TableTest({ + "scenario | input | expected ", + "unset disables | | FULLY_DISABLED ", + "empty string disables | '' | FULLY_DISABLED ", + "unparseable value disables | bad | FULLY_DISABLED ", + "explicit false disables | false | FULLY_DISABLED ", + "zero disables | 0 | FULLY_DISABLED ", + "explicit true enables | true | FULLY_ENABLED ", + "one enables | 1 | FULLY_ENABLED ", + "inactive keyword enables inactive | inactive | ENABLED_INACTIVE" + }) + void iastEnabled(String input, ProductActivation expected) { + if (input != null) { + WithConfigExtension.injectSysConfig("iast.enabled", input); + } + + assertEquals(expected, InstrumenterConfig.get().getIastActivation()); + } + + @TableTest({ + "scenario | input | expected", + "unset defaults to false | | false ", + "explicit false disables | false | false ", + "explicit true enables | true | true ", + "one enables | 1 | true ", + "zero disables | 0 | false " + }) + void dataStreamsEnabled(String input, boolean expected) { + if (input != null) { + WithConfigExtension.injectSysConfig("data.streams.enabled", input); + } + + assertEquals(expected, InstrumenterConfig.get().isDataStreamsEnabled()); + } + + @Test + void dataStreamsEnabledDefaultsToFalse() { + assertFalse(InstrumenterConfig.get().isDataStreamsEnabled()); + } + + @Test + @WithConfig(key = "DATA_STREAMS_ENABLED", value = "true", env = true) + void dataStreamsEnabledViaEnvVar() { + assertTrue(InstrumenterConfig.get().isDataStreamsEnabled()); + } +} diff --git a/internal-api/src/test/java/datadog/trace/api/sampling/SamplingMechanismTest.java b/internal-api/src/test/java/datadog/trace/api/sampling/SamplingMechanismTest.java new file mode 100644 index 00000000000..cd07af4fe78 --- /dev/null +++ b/internal-api/src/test/java/datadog/trace/api/sampling/SamplingMechanismTest.java @@ -0,0 +1,174 @@ +package datadog.trace.api.sampling; + +import static datadog.trace.api.config.GeneralConfig.APM_TRACING_ENABLED; +import static datadog.trace.api.config.GeneralConfig.DATA_STREAMS_ENABLED; +import static datadog.trace.api.sampling.PrioritySampling.SAMPLER_DROP; +import static datadog.trace.api.sampling.PrioritySampling.SAMPLER_KEEP; +import static datadog.trace.api.sampling.PrioritySampling.UNSET; +import static datadog.trace.api.sampling.PrioritySampling.USER_DROP; +import static datadog.trace.api.sampling.PrioritySampling.USER_KEEP; +import static datadog.trace.api.sampling.SamplingMechanism.AGENT_RATE; +import static datadog.trace.api.sampling.SamplingMechanism.APPSEC; +import static datadog.trace.api.sampling.SamplingMechanism.DATA_JOBS; +import static datadog.trace.api.sampling.SamplingMechanism.DATA_STREAMS; +import static datadog.trace.api.sampling.SamplingMechanism.DEFAULT; +import static datadog.trace.api.sampling.SamplingMechanism.EXTERNAL_OVERRIDE; +import static datadog.trace.api.sampling.SamplingMechanism.LOCAL_USER_RULE; +import static datadog.trace.api.sampling.SamplingMechanism.MANUAL; +import static datadog.trace.api.sampling.SamplingMechanism.REMOTE_AUTO_RATE; +import static datadog.trace.api.sampling.SamplingMechanism.REMOTE_USER_RATE; +import static datadog.trace.api.sampling.SamplingMechanism.UNKNOWN; +import static datadog.trace.api.sampling.SamplingMechanism.canAvoidSamplingPriorityLock; +import static datadog.trace.api.sampling.SamplingMechanism.validateWithSamplingPriority; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import datadog.trace.test.junit.utils.config.WithConfig; +import datadog.trace.test.junit.utils.config.WithConfigExtension; +import java.util.stream.Stream; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.Arguments; +import org.junit.jupiter.params.provider.MethodSource; + +@ExtendWith(WithConfigExtension.class) +class SamplingMechanismTest { + + // one below USER_DROP / one above USER_KEEP: neither a valid sampler nor user priority value + private static final byte USER_DROP_X = (byte) (USER_DROP - 1); + private static final byte USER_KEEP_X = (byte) (USER_KEEP + 1); + + @ParameterizedTest + @MethodSource("testValidationArguments") + void testValidation(byte mechanism, byte priority, boolean valid) { + assertEquals(valid, validateWithSamplingPriority(mechanism, priority)); + } + + private static Stream testValidationArguments() { + return Stream.of( + Arguments.of(UNKNOWN, UNSET, true), + Arguments.of(UNKNOWN, SAMPLER_DROP, true), + Arguments.of(UNKNOWN, SAMPLER_KEEP, true), + Arguments.of(UNKNOWN, USER_DROP, true), + Arguments.of(UNKNOWN, USER_KEEP, true), + Arguments.of(UNKNOWN, USER_DROP_X, true), + Arguments.of(UNKNOWN, USER_KEEP_X, true), + Arguments.of(DEFAULT, UNSET, false), + Arguments.of(DEFAULT, SAMPLER_DROP, true), + Arguments.of(DEFAULT, SAMPLER_KEEP, true), + Arguments.of(DEFAULT, USER_DROP, false), + Arguments.of(DEFAULT, USER_KEEP, false), + Arguments.of(DEFAULT, USER_DROP_X, false), + Arguments.of(DEFAULT, USER_KEEP_X, false), + Arguments.of(AGENT_RATE, UNSET, false), + Arguments.of(AGENT_RATE, SAMPLER_DROP, true), + Arguments.of(AGENT_RATE, SAMPLER_KEEP, true), + Arguments.of(AGENT_RATE, USER_DROP, false), + Arguments.of(AGENT_RATE, USER_KEEP, false), + Arguments.of(AGENT_RATE, USER_DROP_X, false), + Arguments.of(AGENT_RATE, USER_KEEP_X, false), + Arguments.of(REMOTE_AUTO_RATE, UNSET, false), + Arguments.of(REMOTE_AUTO_RATE, SAMPLER_DROP, true), + Arguments.of(REMOTE_AUTO_RATE, SAMPLER_KEEP, true), + Arguments.of(REMOTE_AUTO_RATE, USER_DROP, false), + Arguments.of(REMOTE_AUTO_RATE, USER_KEEP, false), + Arguments.of(REMOTE_AUTO_RATE, USER_DROP_X, false), + Arguments.of(REMOTE_AUTO_RATE, USER_KEEP_X, false), + Arguments.of(LOCAL_USER_RULE, UNSET, false), + Arguments.of(LOCAL_USER_RULE, SAMPLER_DROP, false), + Arguments.of(LOCAL_USER_RULE, SAMPLER_KEEP, false), + Arguments.of(LOCAL_USER_RULE, USER_DROP, true), + Arguments.of(LOCAL_USER_RULE, USER_KEEP, true), + Arguments.of(LOCAL_USER_RULE, USER_DROP_X, false), + Arguments.of(LOCAL_USER_RULE, USER_KEEP_X, false), + Arguments.of(MANUAL, UNSET, false), + Arguments.of(MANUAL, SAMPLER_DROP, false), + Arguments.of(MANUAL, SAMPLER_KEEP, false), + Arguments.of(MANUAL, USER_DROP, true), + Arguments.of(MANUAL, USER_KEEP, true), + Arguments.of(MANUAL, USER_DROP_X, false), + Arguments.of(MANUAL, USER_KEEP_X, false), + Arguments.of(REMOTE_USER_RATE, UNSET, false), + Arguments.of(REMOTE_USER_RATE, SAMPLER_DROP, false), + Arguments.of(REMOTE_USER_RATE, SAMPLER_KEEP, false), + Arguments.of(REMOTE_USER_RATE, USER_DROP, true), + Arguments.of(REMOTE_USER_RATE, USER_KEEP, true), + Arguments.of(REMOTE_USER_RATE, USER_DROP_X, false), + Arguments.of(REMOTE_USER_RATE, USER_KEEP_X, false), + Arguments.of(APPSEC, UNSET, false), + Arguments.of(APPSEC, SAMPLER_DROP, true), + Arguments.of(APPSEC, SAMPLER_KEEP, true), + Arguments.of(APPSEC, USER_DROP, false), + Arguments.of(APPSEC, USER_KEEP, true), + Arguments.of(APPSEC, USER_DROP_X, false), + Arguments.of(APPSEC, USER_KEEP_X, false), + Arguments.of(DATA_JOBS, UNSET, false), + Arguments.of(DATA_JOBS, SAMPLER_DROP, false), + Arguments.of(DATA_JOBS, SAMPLER_KEEP, false), + Arguments.of(DATA_JOBS, USER_DROP, false), + Arguments.of(DATA_JOBS, USER_KEEP, true), + Arguments.of(DATA_JOBS, USER_DROP_X, false), + Arguments.of(DATA_JOBS, USER_KEEP_X, false), + Arguments.of(DATA_STREAMS, UNSET, false), + Arguments.of(DATA_STREAMS, SAMPLER_DROP, false), + Arguments.of(DATA_STREAMS, SAMPLER_KEEP, false), + Arguments.of(DATA_STREAMS, USER_DROP, true), + Arguments.of(DATA_STREAMS, USER_KEEP, false), + Arguments.of(DATA_STREAMS, USER_DROP_X, false), + Arguments.of(DATA_STREAMS, USER_KEEP_X, false), + Arguments.of(EXTERNAL_OVERRIDE, UNSET, false), + Arguments.of(EXTERNAL_OVERRIDE, SAMPLER_DROP, false), + Arguments.of(EXTERNAL_OVERRIDE, SAMPLER_KEEP, false), + Arguments.of(EXTERNAL_OVERRIDE, USER_DROP, false), + Arguments.of(EXTERNAL_OVERRIDE, USER_KEEP, false), + Arguments.of(EXTERNAL_OVERRIDE, USER_DROP_X, false), + Arguments.of(EXTERNAL_OVERRIDE, USER_KEEP_X, false)); + } + + @ParameterizedTest + @MethodSource("testCanAvoidSamplingPriorityLockArguments") + @WithConfig(key = APM_TRACING_ENABLED, value = "false") + void testCanAvoidSamplingPriorityLock(byte mechanism, byte priority, boolean valid) { + assertEquals(valid, canAvoidSamplingPriorityLock(priority, mechanism)); + } + + private static Stream testCanAvoidSamplingPriorityLockArguments() { + return Stream.of( + Arguments.of(APPSEC, UNSET, true), + Arguments.of(APPSEC, SAMPLER_KEEP, true), + Arguments.of(UNKNOWN, SAMPLER_KEEP, false), + Arguments.of(DEFAULT, SAMPLER_KEEP, false), + Arguments.of(AGENT_RATE, SAMPLER_KEEP, false), + Arguments.of(REMOTE_AUTO_RATE, SAMPLER_KEEP, false), + Arguments.of(LOCAL_USER_RULE, SAMPLER_KEEP, false), + Arguments.of(MANUAL, SAMPLER_KEEP, false), + Arguments.of(REMOTE_USER_RATE, SAMPLER_KEEP, false), + Arguments.of(DATA_JOBS, SAMPLER_KEEP, false), + // DSM is left at its config default (disabled) for this parameterized run, so the + // DATA_STREAMS case is false here for a config reason. The two dedicated tests below + // cover the mechanism itself with data.streams.enabled explicitly set both ways. + Arguments.of(DATA_STREAMS, SAMPLER_KEEP, false), + Arguments.of(EXTERNAL_OVERRIDE, SAMPLER_KEEP, false)); + } + + @Test + @WithConfig(key = DATA_STREAMS_ENABLED, value = "true") + void dataStreamsMechanismCanAvoidSamplingPriorityLockWhenDataStreamsEnabled() { + // The DATA_STREAMS case is priority-independent: the mechanism alone unlocks the priority. + assertTrue(canAvoidSamplingPriorityLock(USER_DROP, DATA_STREAMS)); + assertTrue(canAvoidSamplingPriorityLock(SAMPLER_KEEP, DATA_STREAMS)); + assertTrue(canAvoidSamplingPriorityLock(UNSET, DATA_STREAMS)); + // Enabling DSM must not unlock any other mechanism. + assertFalse(canAvoidSamplingPriorityLock(USER_DROP, MANUAL)); + assertFalse(canAvoidSamplingPriorityLock(USER_DROP, DEFAULT)); + } + + @Test + @WithConfig(key = DATA_STREAMS_ENABLED, value = "false") + void dataStreamsMechanismCannotAvoidSamplingPriorityLockWhenDataStreamsDisabled() { + assertFalse(canAvoidSamplingPriorityLock(USER_DROP, DATA_STREAMS)); + assertFalse(canAvoidSamplingPriorityLock(SAMPLER_KEEP, DATA_STREAMS)); + } +}