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..7b1aa90cd82 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 @@ -20,6 +20,8 @@ import datadog.trace.agent.tooling.Instrumenter; import datadog.trace.agent.tooling.InstrumenterModule; import datadog.trace.api.Config; +import datadog.trace.api.sampling.PrioritySampling; +import datadog.trace.api.sampling.SamplingMechanism; import datadog.trace.bootstrap.InstrumentationContext; import datadog.trace.bootstrap.instrumentation.api.AgentScope; import datadog.trace.bootstrap.instrumentation.api.AgentSpan; @@ -44,11 +46,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 @@ -254,6 +256,13 @@ public static AgentScope onEnter(@Advice.This KafkaConsumer consumer) { if (traceConfig().isDataStreamsEnabled()) { final AgentSpan span = startSpan(JAVA_KAFKA.toString(), KAFKA_POLL); + // setSamplingPriority is trace-level: it resolves to the local root span. This 2-arg + // startSpan honours the active scope, so `span` may be a child of a customer trace. Only + // force the DSM-only drop when `span` is the local root, otherwise we would silently drop + // that whole customer trace. + if (!KafkaDecorator.TRACING_ENABLED && span.getLocalRootSpan() == span) { + span.setSamplingPriority(PrioritySampling.USER_DROP, SamplingMechanism.DATA_STREAMS); + } return activateSpan(span); } return null; 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..fe232fea06d 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,6 +12,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_clients.KafkaDecorator.JAVA_KAFKA; import static datadog.trace.instrumentation.kafka_clients.KafkaDecorator.KAFKA_PRODUCE; import static datadog.trace.instrumentation.kafka_clients.KafkaDecorator.PRODUCER_DECORATE; @@ -37,6 +38,8 @@ import datadog.trace.api.datastreams.DataStreamsTags; import datadog.trace.api.datastreams.DataStreamsTransactionExtractor; import datadog.trace.api.datastreams.StatsPoint; +import datadog.trace.api.sampling.PrioritySampling; +import datadog.trace.api.sampling.SamplingMechanism; import datadog.trace.bootstrap.InstrumentationContext; import datadog.trace.bootstrap.instrumentation.api.AgentScope; import datadog.trace.bootstrap.instrumentation.api.AgentSpan; @@ -58,11 +61,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 @@ -165,6 +168,15 @@ public static AgentScope onEnter( } else { span = startSpan(JAVA_KAFKA.toString(), KAFKA_PRODUCE); callbackParentSpan = localActiveSpan; + // setSamplingPriority is trace-level: it resolves to the local root span. This 2-arg + // startSpan honours the active scope, so `span` may be a child of a customer trace + // (localActiveSpan above). Only force the DSM-only drop when `span` is the local root, + // otherwise we would silently drop that whole customer trace. + if (!KafkaDecorator.TRACING_ENABLED + && traceConfig().isDataStreamsEnabled() + && span.getLocalRootSpan() == span) { + span.setSamplingPriority(PrioritySampling.USER_DROP, SamplingMechanism.DATA_STREAMS); + } } PRODUCER_DECORATE.afterStart(span); PRODUCER_DECORATE.onProduce(span, record, producerConfig, 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..cd242c152ec 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 @@ -28,6 +28,8 @@ import datadog.trace.api.datastreams.DataStreamsContext; import datadog.trace.api.datastreams.DataStreamsTags; import datadog.trace.api.datastreams.DataStreamsTransactionExtractor; +import datadog.trace.api.sampling.PrioritySampling; +import datadog.trace.api.sampling.SamplingMechanism; import datadog.trace.bootstrap.instrumentation.api.AgentSpan; import datadog.trace.bootstrap.instrumentation.api.AgentSpanContext; import datadog.trace.bootstrap.instrumentation.api.AgentTracer; @@ -120,6 +122,11 @@ protected void startNewRecordSpan(ConsumerRecord val) { // spans are written out together by TraceStructureWriter when running in strict mode } + if (spanContext == null + && !KafkaDecorator.TRACING_ENABLED + && traceConfig().isDataStreamsEnabled()) { + span.setSamplingPriority(PrioritySampling.USER_DROP, SamplingMechanism.DATA_STREAMS); + } DataStreamsTags tags = create("kafka", INBOUND, val.topic(), group, clusterId); final long payloadSize = traceConfig().isDataStreamsEnabled() ? computePayloadSizeBytes(val) : 0; @@ -141,6 +148,9 @@ protected void startNewRecordSpan(ConsumerRecord val) { } } else { span = startSpan(JAVA_KAFKA.toString(), operationName, null); + if (!KafkaDecorator.TRACING_ENABLED && traceConfig().isDataStreamsEnabled()) { + span.setSamplingPriority(PrioritySampling.USER_DROP, SamplingMechanism.DATA_STREAMS); + } } if (val.value() == null) { span.setTag(InstrumentationTags.TOMBSTONE, true); 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..7bf4553d867 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 @@ -14,6 +14,7 @@ import datadog.trace.agent.test.asserts.TraceAssert import datadog.trace.agent.test.naming.VersionedNamingTestBase import datadog.trace.api.Config import datadog.trace.api.DDTags +import datadog.trace.api.sampling.PrioritySampling import datadog.trace.bootstrap.instrumentation.api.InstrumentationTags import datadog.trace.bootstrap.instrumentation.api.Tags import datadog.trace.common.writer.ListWriter @@ -1544,3 +1545,345 @@ class KafkaClientBadBase64HeaderForkedTest extends InstrumentationSpecification producer?.close() } } + +// DSM billing-suppression coverage: when tracing is disabled for kafka (globally, via +// integrations.enabled) but DSM is enabled, a genuinely local-root produce/consume span (no +// extracted/real parent trace context) must have its sampling priority forced to USER_DROP so it +// does not count towards APM billing, while a span that joins a real propagated trace must not be. +class KafkaClientDataStreamsOnlyLocalRootForkedTest extends KafkaClientTestBase { + @Override + void configurePreAgent() { + super.configurePreAgent() + injectSysConfig("integrations.enabled", "false") + injectSysConfig("data.streams.enabled", "true") + } + + @Override + String service() { + return "kafka" + } + + @Override + boolean hasQueueSpan() { + return false + } + + @Override + boolean splitByDestination() { + return false + } + + @Override + boolean isDataStreamsEnabled() { + return true + } + + def "local-root produce and consume spans are forced to USER_DROP 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 senderProps = KafkaTestUtils.senderProps(embeddedKafka.getBrokersAsString()) + def producer = new KafkaProducer<>(senderProps, new StringSerializer(), new StringSerializer()) + 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" + def record = new ProducerRecord(SHARED_TOPIC, kafkaPartition, null, "local-root-message") + producer.send(record).get() + + then: "the produce span's trace is forced to USER_DROP to suppress APM billing" + TEST_WRITER.waitForTraces(1) + def producedSpan = TEST_WRITER[0][0] + producedSpan.getSamplingPriority() == PrioritySampling.USER_DROP + + when: "the message is consumed" + def pollResult = KafkaTestUtils.getRecords(consumer) + def recs = pollResult.records(new TopicPartition(SHARED_TOPIC, kafkaPartition)).iterator() + + then: "the consume span's trace is also forced to USER_DROP" + recs.hasNext() + recs.next().value() == "local-root-message" + !recs.hasNext() + TEST_WRITER.waitForTraces(2) + def consumedSpan = TEST_WRITER[1][0] + consumedSpan.getSamplingPriority() == PrioritySampling.USER_DROP + + cleanup: + consumer?.close() + producer?.close() + } +} + +// Regression guard: a span that joins a real, externally-propagated Datadog trace (extracted +// x-datadog-trace-id/x-datadog-parent-id headers) must NOT be forced to USER_DROP, even under the +// same "kafka tracing disabled + DSM enabled" configuration as above, since it is not a local root. +class KafkaClientDataStreamsOnlyExtractedParentForkedTest extends KafkaClientTestBase { + @Override + void configurePreAgent() { + super.configurePreAgent() + injectSysConfig("integrations.enabled", "false") + injectSysConfig("data.streams.enabled", "true") + } + + @Override + String service() { + return "kafka" + } + + @Override + boolean hasQueueSpan() { + return false + } + + @Override + boolean splitByDestination() { + return false + } + + @Override + boolean isDataStreamsEnabled() { + return true + } + + def "produce and consume spans with an extracted parent trace context are not forced to USER_DROP"() { + 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 senderProps = KafkaTestUtils.senderProps(embeddedKafka.getBrokersAsString()) + def producer = new KafkaProducer<>(senderProps, new StringSerializer(), new StringSerializer()) + 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" + def record = new ProducerRecord(SHARED_TOPIC, kafkaPartition, null, "propagated-trace-message", headers) + producer.send(record).get() + + then: "the produce span joins the propagated trace and is NOT forced to USER_DROP" + TEST_WRITER.waitForTraces(1) + def producedSpan = TEST_WRITER[0][0] + producedSpan.traceId.toLong() == existingTraceId + producedSpan.parentId == existingSpanId + // The injected headers carry no x-datadog-sampling-priority, so the extracted context's + // priority is UNSET and the normal sampler decides. Assert the concrete resulting value + // rather than just "not USER_DROP", so any regression that lands on a different-but-also + // wrong priority is caught too. + producedSpan.getSamplingPriority() == PrioritySampling.SAMPLER_KEEP + + when: "the message is consumed" + def pollResult = KafkaTestUtils.getRecords(consumer) + def recs = pollResult.records(new TopicPartition(SHARED_TOPIC, kafkaPartition)).iterator() + + then: "the consume span also joins the propagated trace and is NOT forced to USER_DROP" + recs.hasNext() + recs.next().value() == "propagated-trace-message" + !recs.hasNext() + TEST_WRITER.waitForTraces(2) + def consumedSpan = TEST_WRITER[1][0] + consumedSpan.traceId.toLong() == existingTraceId + consumedSpan.getSamplingPriority() == PrioritySampling.SAMPLER_KEEP + + cleanup: + consumer?.close() + producer?.close() + } +} + +// Confirms the per-integration override (trace.kafka.enabled=false) suppresses billing identically +// to the global integrations.enabled=false toggle used above, for a genuine local-root span. +class KafkaClientDataStreamsOnlyIntegrationOverrideForkedTest extends KafkaClientTestBase { + @Override + void configurePreAgent() { + super.configurePreAgent() + injectSysConfig("trace.kafka.enabled", "false") + injectSysConfig("data.streams.enabled", "true") + } + + @Override + String service() { + return "kafka" + } + + @Override + boolean hasQueueSpan() { + return false + } + + @Override + boolean splitByDestination() { + return false + } + + @Override + boolean isDataStreamsEnabled() { + return true + } + + def "local-root produce and consume spans are forced to USER_DROP 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 senderProps = KafkaTestUtils.senderProps(embeddedKafka.getBrokersAsString()) + def producer = new KafkaProducer<>(senderProps, new StringSerializer(), new StringSerializer()) + 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" + def record = new ProducerRecord(SHARED_TOPIC, kafkaPartition, null, "local-root-message") + producer.send(record).get() + + then: "the produce span's trace is forced to USER_DROP to suppress APM billing" + TEST_WRITER.waitForTraces(1) + def producedSpan = TEST_WRITER[0][0] + producedSpan.getSamplingPriority() == PrioritySampling.USER_DROP + + when: "the message is consumed" + def pollResult = KafkaTestUtils.getRecords(consumer) + def recs = pollResult.records(new TopicPartition(SHARED_TOPIC, kafkaPartition)).iterator() + + then: "the consume span's trace is also forced to USER_DROP" + recs.hasNext() + recs.next().value() == "local-root-message" + !recs.hasNext() + TEST_WRITER.waitForTraces(2) + def consumedSpan = TEST_WRITER[1][0] + consumedSpan.getSamplingPriority() == PrioritySampling.USER_DROP + + cleanup: + consumer?.close() + producer?.close() + } +} + +// Regression guard for the producer suppression site: producing a message from inside an already +// active local trace (e.g. an instrumented HTTP request) must NOT force that customer trace to +// USER_DROP. The produce span is created with a scope-honouring startSpan overload, so it becomes +// a child of the active span, and setSamplingPriority is trace-level - without the local-root +// check on the suppression guard the whole surrounding trace would be silently dropped. +class KafkaClientDataStreamsOnlyActiveLocalTraceForkedTest extends KafkaClientTestBase { + @Override + void configurePreAgent() { + super.configurePreAgent() + injectSysConfig("integrations.enabled", "false") + injectSysConfig("data.streams.enabled", "true") + } + + @Override + String service() { + return "kafka" + } + + @Override + boolean hasQueueSpan() { + return false + } + + @Override + boolean splitByDestination() { + return false + } + + @Override + boolean isDataStreamsEnabled() { + return true + } + + def "producing inside an active local trace does not force that trace to USER_DROP"() { + setup: + def kafkaPartition = 0 + def senderProps = KafkaTestUtils.senderProps(embeddedKafka.getBrokersAsString()) + def producer = new KafkaProducer<>(senderProps, new StringSerializer(), new StringSerializer()) + + when: "a message is produced from within an already active local trace" + runUnderTrace("parent") { + producer.send(new ProducerRecord(SHARED_TOPIC, kafkaPartition, null, "in-active-trace")).get() + } + + then: "the surrounding customer trace is not force-dropped" + TEST_WRITER.waitForTraces(1) + def trace = TEST_WRITER[0] + def localRoot = trace[0].localRootSpan + localRoot.operationName.toString() == "parent" + localRoot.getSamplingPriority() != PrioritySampling.USER_DROP + + cleanup: + producer?.close() + } +} + +// Regression guard for the poll-span suppression site: KafkaConsumerInfoInstrumentation's +// RecordsAdvice creates a standalone "kafka.poll" span/trace around every consumer.poll() call +// whenever DSM is enabled, regardless of whether Kafka APM tracing itself is enabled. Unlike the +// other cases above, this test intentionally keeps the "kafka.poll" trace instead of relying on +// the base class's DROP_KAFKA_POLL filter, so the assertion actually observes what gets written. +class KafkaClientDataStreamsOnlyPollSpanForkedTest extends KafkaClientTestBase { + 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") + injectSysConfig("data.streams.enabled", "true") + } + + def setup() { + // Undo the base class's DROP_KAFKA_POLL filter: this test needs to see the "kafka.poll" trace. + TEST_WRITER.setFilter(ACCEPT_ALL) + } + + @Override + String service() { + return "kafka" + } + + @Override + boolean hasQueueSpan() { + return false + } + + @Override + boolean splitByDestination() { + return false + } + + @Override + boolean isDataStreamsEnabled() { + return true + } + + def "poll span is forced to USER_DROP when kafka tracing is disabled and DSM is enabled"() { + 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: "the standalone kafka.poll trace is forced to USER_DROP, so it is not billed as APM" + TEST_WRITER.waitForTraces(1) + def trace = TEST_WRITER[0] + trace.size() == 1 + trace[0].getResourceName().toString() == "kafka.poll" + trace[0].getSamplingPriority() == PrioritySampling.USER_DROP + + cleanup: + consumer?.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..5a1c2f05989 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,10 +4,13 @@ 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; +import datadog.trace.api.sampling.PrioritySampling; +import datadog.trace.api.sampling.SamplingMechanism; import datadog.trace.bootstrap.InstrumentationContext; import datadog.trace.bootstrap.instrumentation.api.AgentScope; import datadog.trace.bootstrap.instrumentation.api.AgentSpan; @@ -55,6 +58,15 @@ public static AgentScope onEnter( } else { span = startSpan(JAVA_KAFKA.toString(), KAFKA_PRODUCE); callbackParentSpan = localActiveSpan; + // setSamplingPriority is trace-level: it resolves to the local root span. This 2-arg + // startSpan honours the active scope, so `span` may be a child of a customer trace + // (localActiveSpan above). Only force the DSM-only drop when `span` is the local root, + // otherwise we would silently drop that whole customer trace. + if (!KafkaDecorator.TRACING_ENABLED + && traceConfig().isDataStreamsEnabled() + && span.getLocalRootSpan() == span) { + span.setSamplingPriority(PrioritySampling.USER_DROP, SamplingMechanism.DATA_STREAMS); + } } PRODUCER_DECORATE.afterStart(span); PRODUCER_DECORATE.onProduce(span, record, producerConfig, 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..a3c6b964c24 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 @@ -8,6 +8,8 @@ import static datadog.trace.instrumentation.kafka_clients38.KafkaDecorator.KAFKA_POLL; import datadog.trace.api.Config; +import datadog.trace.api.sampling.PrioritySampling; +import datadog.trace.api.sampling.SamplingMechanism; import datadog.trace.bootstrap.InstrumentationContext; import datadog.trace.bootstrap.instrumentation.api.AgentScope; import datadog.trace.bootstrap.instrumentation.api.AgentSpan; @@ -41,6 +43,13 @@ public static AgentScope onEnter(@Advice.This ConsumerDelegate consumer) { if (traceConfig().isDataStreamsEnabled()) { final AgentSpan span = startSpan(JAVA_KAFKA.toString(), KAFKA_POLL); + // setSamplingPriority is trace-level: it resolves to the local root span. This 2-arg + // startSpan honours the active scope, so `span` may be a child of a customer trace. Only + // force the DSM-only drop when `span` is the local root, otherwise we would silently drop + // that whole customer trace. + if (!KafkaDecorator.TRACING_ENABLED && span.getLocalRootSpan() == span) { + span.setSamplingPriority(PrioritySampling.USER_DROP, SamplingMechanism.DATA_STREAMS); + } return activateSpan(span); } return null; 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..a14c3da29d6 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 @@ -22,6 +22,8 @@ import datadog.trace.api.datastreams.DataStreamsContext; import datadog.trace.api.datastreams.DataStreamsTags; import datadog.trace.api.datastreams.DataStreamsTransactionExtractor; +import datadog.trace.api.sampling.PrioritySampling; +import datadog.trace.api.sampling.SamplingMechanism; import datadog.trace.bootstrap.instrumentation.api.AgentSpan; import datadog.trace.bootstrap.instrumentation.api.AgentSpanContext; import datadog.trace.bootstrap.instrumentation.api.AgentTracer; @@ -116,6 +118,11 @@ protected void startNewRecordSpan(ConsumerRecord val) { // spans are written out together by TraceStructureWriter when running in strict mode } + if (spanContext == null + && !KafkaDecorator.TRACING_ENABLED + && traceConfig().isDataStreamsEnabled()) { + span.setSamplingPriority(PrioritySampling.USER_DROP, SamplingMechanism.DATA_STREAMS); + } DataStreamsTags tags = create("kafka", INBOUND, val.topic(), group, clusterId); final long payloadSize = traceConfig().isDataStreamsEnabled() ? Utils.computePayloadSizeBytes(val) : 0; @@ -137,6 +144,9 @@ protected void startNewRecordSpan(ConsumerRecord val) { } } else { span = startSpan(JAVA_KAFKA.toString(), operationName, null); + if (!KafkaDecorator.TRACING_ENABLED && traceConfig().isDataStreamsEnabled()) { + span.setSamplingPriority(PrioritySampling.USER_DROP, SamplingMechanism.DATA_STREAMS); + } } if (val.value() == null) { span.setTag(InstrumentationTags.TOMBSTONE, true); 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..98c47a98e37 --- /dev/null +++ b/dd-java-agent/instrumentation/kafka/kafka-clients-3.8/src/test/groovy/KafkaClientDataStreamsOnlyForkedTest.groovy @@ -0,0 +1,168 @@ +import datadog.trace.agent.test.InstrumentationSpecification +import datadog.trace.api.sampling.PrioritySampling +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 billing-suppression coverage for kafka-clients-3.8, mirroring the kafka-clients-0.11 suite. + * + *

    Both scenarios run under "kafka tracing disabled (integrations.enabled=false) + DSM enabled", + * the configuration in which the suppression guard is active. These specs deliberately do not + * extend {@code KafkaClientTestBase}: that base asserts Code Origin tags, 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 + void configurePreAgent() { + super.configurePreAgent() + injectSysConfig("integrations.enabled", "false") + injectSysConfig("data.streams.enabled", "true") + } + + /** + * Locates a written span by operation name. The consumer's kafka.poll spans interleave + * unpredictably with the produce/consume traces, so indexing into TEST_WRITER is unreliable. + */ + protected findSpan(String operationName) { + for (int i = 0; i < 100; i++) { + def span = TEST_WRITER.flatten().find { it.operationName.toString() == operationName } + if (span != null) { + return span + } + Thread.sleep(100) + } + return null + } + + protected KafkaProducer newProducer() { + return new KafkaProducer( + KafkaTestUtils.producerProps(embeddedKafka.getBrokersAsString()), + new StringSerializer(), + new StringSerializer()) + } +} + +/** + * A genuinely local-root produce/consume span must have its sampling priority forced to USER_DROP + * so it does not count towards APM billing. + */ +class KafkaClientDataStreamsOnlyLocalRootForkedTest extends KafkaClientDataStreamsOnlyForkedTest { + + def "local-root produce and consume spans are forced to USER_DROP"() { + 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() + + then: "the produce span's trace is forced to USER_DROP to suppress APM billing" + TEST_WRITER.waitForTraces(1) + def produceSpan = findSpan("kafka.produce") + produceSpan != null + produceSpan.getSamplingPriority() == PrioritySampling.USER_DROP + + when: "the message is consumed" + def recs = KafkaTestUtils.getRecords(consumer) + .records(new TopicPartition(SHARED_TOPIC, kafkaPartition)).iterator() + + then: "the consume span's trace is also forced to USER_DROP" + recs.hasNext() + recs.next().value() == "local-root-message" + !recs.hasNext() + TEST_WRITER.waitForTraces(2) + def consumeSpan = findSpan("kafka.consume") + consumeSpan != null + consumeSpan.getSamplingPriority() == PrioritySampling.USER_DROP + + cleanup: + consumer?.close() + producer?.close() + } +} + +/** + * Regression guard for the producer suppression site: producing a message from inside an already + * active local trace must NOT force that customer trace to USER_DROP. The produce span is created + * with a scope-honouring startSpan overload, so it becomes a child of the active span, and + * setSamplingPriority is trace-level. + */ +class KafkaClientDataStreamsOnlyActiveLocalTraceForkedTest extends KafkaClientDataStreamsOnlyForkedTest { + + def "producing inside an active local trace does not force that trace to USER_DROP"() { + 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 is not force-dropped" + TEST_WRITER.waitForTraces(1) + def localRoot = TEST_WRITER[0][0].localRootSpan + localRoot.operationName.toString() == "parent" + localRoot.getSamplingPriority() != PrioritySampling.USER_DROP + + cleanup: + producer?.close() + } +} + +/** + * Regression guard for the poll-span suppression site: KafkaConsumerInfoInstrumentation's + * RecordsAdvice creates a standalone "kafka.poll" span/trace around every consumer.poll() call + * whenever DSM is enabled, regardless of whether Kafka APM tracing itself is enabled. + */ +class KafkaClientDataStreamsOnlyPollSpanForkedTest extends KafkaClientDataStreamsOnlyForkedTest { + + def "poll span is forced to USER_DROP"() { + 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: "the standalone kafka.poll trace is forced to USER_DROP, so it is not billed as APM" + def pollSpan = findSpan("kafka.poll") + pollSpan != null + pollSpan.getSamplingPriority() == PrioritySampling.USER_DROP + + cleanup: + consumer?.close() + } +} 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..e2842543d98 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,6 +8,7 @@ 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; @@ -41,6 +42,8 @@ import datadog.trace.api.Config; import datadog.trace.api.datastreams.DataStreamsContext; import datadog.trace.api.datastreams.DataStreamsTags; +import datadog.trace.api.sampling.PrioritySampling; +import datadog.trace.api.sampling.SamplingMechanism; import datadog.trace.bootstrap.InstrumentationContext; import datadog.trace.bootstrap.instrumentation.api.AgentScope; import datadog.trace.bootstrap.instrumentation.api.AgentSpan; @@ -57,11 +60,25 @@ 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); + } + + // setSamplingPriority is trace-level: it resolves to the local root span. Only force the + // DSM-only drop when this instrumentation owns the whole local trace, i.e. nothing was + // active when we started (no local parent, no header-extracted parent) and the local root + // really is the first span we created here. + private static void maybeDropForDataStreamsOnly( + final AgentSpan span, final AgentSpan localActiveSpan, final AgentSpan ourLocalRoot) { + if (!KafkaStreamsDecorator.TRACING_ENABLED + && traceConfig().isDataStreamsEnabled() + && localActiveSpan == null + && span.getLocalRootSpan() == ourLocalRoot) { + span.setSamplingPriority(PrioritySampling.USER_DROP, SamplingMechanism.DATA_STREAMS); + } } @Override @@ -276,6 +293,12 @@ public static void start( return; } + // Captured before any span is created. A non-null value here means this record is being + // consumed either inside a locally active trace, or under a context that the sibling + // ContextPropagationAdvice extracted from the record headers and attached to the scope + // (it is registered on the same method and runs first). Either way the spans created + // below will not own the resulting trace. + final AgentSpan localActiveSpan = activeSpan(); AgentSpan span, queueSpan = null; StreamTaskContext streamTaskContext = InstrumentationContext.get(StreamTask.class, StreamTaskContext.class).get(task); @@ -294,6 +317,8 @@ public static void start( // spans are written out together by TraceStructureWriter when running in strict mode } + final AgentSpan ourLocalRoot = queueSpan == null ? span : queueSpan; + maybeDropForDataStreamsOnly(span, localActiveSpan, ourLocalRoot); String applicationId = null; if (streamTaskContext != null) { applicationId = streamTaskContext.getApplicationId(); @@ -342,6 +367,12 @@ public static void start( return; } + // Captured before any span is created. A non-null value here means this record is being + // consumed either inside a locally active trace, or under a context that the sibling + // ContextPropagationAdvice extracted from the record headers and attached to the scope + // (it is registered on the same method and runs first). Either way the spans created + // below will not own the resulting trace. + final AgentSpan localActiveSpan = activeSpan(); AgentSpan span, queueSpan = null; StreamTaskContext streamTaskContext = InstrumentationContext.get(StreamTask.class, StreamTaskContext.class).get(task); @@ -360,6 +391,8 @@ public static void start( // spans are written out together by TraceStructureWriter when running in strict mode } + final AgentSpan ourLocalRoot = queueSpan == null ? span : queueSpan; + maybeDropForDataStreamsOnly(span, localActiveSpan, ourLocalRoot); String applicationId = null; if (streamTaskContext != null) { applicationId = streamTaskContext.getApplicationId(); 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..185121d447b --- /dev/null +++ b/dd-java-agent/instrumentation/kafka/kafka-streams-0.11/src/test/groovy/KafkaStreamsDataStreamsOnlyForkedTest.groovy @@ -0,0 +1,173 @@ +import datadog.trace.agent.test.InstrumentationSpecification +import datadog.trace.api.config.TraceInstrumentationConfig +import datadog.trace.api.sampling.PrioritySampling +import datadog.trace.bootstrap.instrumentation.api.Tags +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 billing-suppression coverage for the kafka-streams StreamTask consume spans. + * + *

    Both scenarios run under "kafka tracing disabled (integrations.enabled=false) + DSM enabled", + * the configuration in which the suppression guard is active. + */ +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") + injectSysConfig("data.streams.enabled", "true") + } + + @Override + boolean useStrictTraceWrites() { + return false + } + + 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()) + } + + /** + * Polls the test writer until a kafka-streams consume span shows up, so the assertions do not + * depend on how many other traces (produce, poll, downstream produce) are flushed first. + */ + protected findStreamsConsumeSpan() { + for (int i = 0; i < 100; i++) { + def span = TEST_WRITER.flatten().find { + it.operationName.toString() == "kafka.consume" && + it.getTag(Tags.COMPONENT)?.toString() == "java-kafka-streams" + } + if (span != null) { + return span + } + Thread.sleep(100) + } + return null + } +} + +/** + * A record with no propagated Datadog trace context produces a genuinely local-root streams + * consume span, which must be forced to USER_DROP so it does not count towards APM billing. + */ +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 streams consume span is forced to USER_DROP"() { + setup: + def streams = startLowercasingTopology() + def producer = newProducer() + + when: + producer.send(new ProducerRecord(STREAM_PENDING, "LOCAL ROOT")).get() + + then: + def consumeSpan = findStreamsConsumeSpan() + consumeSpan != null + consumeSpan.getSamplingPriority() == PrioritySampling.USER_DROP + + cleanup: + producer?.close() + streams?.close() + } +} + +/** + * Regression guard: a record carrying a real, externally-propagated Datadog trace context must NOT + * have its trace force-dropped. The sibling ContextPropagationAdvice attaches that extracted + * context to the scope before the span-starting advice runs, and setSamplingPriority is + * trace-level, so without the guard the whole propagated trace would be silently dropped. + */ +class KafkaStreamsDataStreamsOnlyExtractedParentForkedTest extends KafkaStreamsDataStreamsOnlyForkedTest { + + def "a streams consume span continuing a propagated trace is not forced to USER_DROP"() { + 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: + def consumeSpan = findStreamsConsumeSpan() + consumeSpan != null + consumeSpan.traceId.toLong() == existingTraceId + consumeSpan.getSamplingPriority() != PrioritySampling.USER_DROP + + 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)); + } +}