diff --git a/examples/src/main/java/software/amazon/lambda/durable/examples/general/OtelExample.java b/examples/src/main/java/software/amazon/lambda/durable/examples/general/OtelExample.java index 3567501c1..fca3e306b 100644 --- a/examples/src/main/java/software/amazon/lambda/durable/examples/general/OtelExample.java +++ b/examples/src/main/java/software/amazon/lambda/durable/examples/general/OtelExample.java @@ -9,7 +9,6 @@ import software.amazon.lambda.durable.DurableContext; import software.amazon.lambda.durable.DurableHandler; import software.amazon.lambda.durable.examples.types.GreetingRequest; -import software.amazon.lambda.durable.otel.DeterministicIdGenerator; import software.amazon.lambda.durable.otel.OpenTelemetryDurablePlugin; /** @@ -40,12 +39,8 @@ public class OtelExample extends DurableHandler { @Override protected DurableConfig createConfiguration() { - var idGenerator = new DeterministicIdGenerator(); - var tracerProvider = SdkTracerProvider.builder() - .setIdGenerator(idGenerator) - .addSpanProcessor(SimpleSpanProcessor.create(LoggingSpanExporter.create())) - .build(); - var otelPlugin = new OpenTelemetryDurablePlugin(tracerProvider, idGenerator); + var otelPlugin = new OpenTelemetryDurablePlugin( + SdkTracerProvider.builder().addSpanProcessor(SimpleSpanProcessor.create(LoggingSpanExporter.create()))); return DurableConfig.builder().withPlugins(otelPlugin).build(); } diff --git a/otel-plugin/src/main/java/software/amazon/lambda/durable/otel/OpenTelemetryDurablePlugin.java b/otel-plugin/src/main/java/software/amazon/lambda/durable/otel/OpenTelemetryDurablePlugin.java index 0e9e93ade..29806a3e9 100644 --- a/otel-plugin/src/main/java/software/amazon/lambda/durable/otel/OpenTelemetryDurablePlugin.java +++ b/otel-plugin/src/main/java/software/amazon/lambda/durable/otel/OpenTelemetryDurablePlugin.java @@ -11,10 +11,10 @@ import io.opentelemetry.api.trace.TraceFlags; import io.opentelemetry.api.trace.TraceState; import io.opentelemetry.api.trace.Tracer; -import io.opentelemetry.api.trace.TracerProvider; import io.opentelemetry.context.Context; import io.opentelemetry.context.Scope; import io.opentelemetry.sdk.trace.SdkTracerProvider; +import io.opentelemetry.sdk.trace.SdkTracerProviderBuilder; import java.time.Instant; import java.util.concurrent.ConcurrentHashMap; import org.slf4j.Logger; @@ -52,17 +52,15 @@ public class OpenTelemetryDurablePlugin implements DurableExecutionPlugin { private static final Logger logger = LoggerFactory.getLogger(OpenTelemetryDurablePlugin.class); private static final String INSTRUMENTATION_NAME = "aws-durable-execution-sdk-java"; - private final TracerProvider tracerProvider; + private final SdkTracerProvider tracerProvider; private final Tracer tracer; private final DeterministicIdGenerator idGenerator; private final ContextExtractor contextExtractor; - private final double samplingRate; private final boolean enableMdc; // Per-invocation state private volatile Span invocationSpan; private volatile String executionArn; - private volatile boolean sampled = true; // Thread-safe storage for operation spans (keyed by operationId) — open spans that need ending private final ConcurrentHashMap operationSpans = new ConcurrentHashMap<>(); @@ -75,64 +73,45 @@ public class OpenTelemetryDurablePlugin implements DurableExecutionPlugin { private final ConcurrentHashMap operationContexts = new ConcurrentHashMap<>(); /** - * Creates an OTel plugin with default settings: X-Ray context extraction, 100% sampling, MDC enabled. + * Creates an OTel plugin with default settings: X-Ray context extraction, MDC enabled. * - * @param tracerProvider the OTel tracer provider (should use {@link DeterministicIdGenerator}) - * @param idGenerator the deterministic ID generator (same instance configured in the tracer provider) + * @param tracerProviderBuilder the tracer provider builder (ID generator will be overridden) */ - public OpenTelemetryDurablePlugin(TracerProvider tracerProvider, DeterministicIdGenerator idGenerator) { - this(tracerProvider, idGenerator, new XRayContextExtractor(), 1.0, true); + public OpenTelemetryDurablePlugin(SdkTracerProviderBuilder tracerProviderBuilder) { + this(tracerProviderBuilder, new XRayContextExtractor(), true); } /** - * Creates an OTel plugin with a custom context extractor and sampling rate, MDC enabled. + * Creates an OTel plugin with a custom context extractor, MDC enabled. * - * @param tracerProvider the OTel tracer provider - * @param idGenerator the deterministic ID generator + * @param tracerProviderBuilder the tracer provider builder (ID generator will be overridden) * @param contextExtractor extracts parent trace context from the Lambda environment - * @param samplingRate value between 0.0 and 1.0 — fraction of executions to trace */ public OpenTelemetryDurablePlugin( - TracerProvider tracerProvider, - DeterministicIdGenerator idGenerator, - ContextExtractor contextExtractor, - double samplingRate) { - this(tracerProvider, idGenerator, contextExtractor, samplingRate, true); + SdkTracerProviderBuilder tracerProviderBuilder, ContextExtractor contextExtractor) { + this(tracerProviderBuilder, contextExtractor, true); } /** * Creates an OTel plugin with full configuration. * - * @param tracerProvider the OTel tracer provider - * @param idGenerator the deterministic ID generator + *

The plugin internally creates a {@link DeterministicIdGenerator} and sets it on the provided builder before + * building the tracer provider. Customers configure exporters, span processors, and samplers on the builder — the + * plugin handles ID generation. + * + * @param tracerProviderBuilder the tracer provider builder (ID generator will be overridden) * @param contextExtractor extracts parent trace context from the Lambda environment - * @param samplingRate value between 0.0 and 1.0 — fraction of executions to trace * @param enableMdc if true, injects trace_id/span_id into SLF4J MDC for log correlation */ public OpenTelemetryDurablePlugin( - TracerProvider tracerProvider, - DeterministicIdGenerator idGenerator, - ContextExtractor contextExtractor, - double samplingRate, - boolean enableMdc) { - this.tracerProvider = tracerProvider; - this.idGenerator = idGenerator; + SdkTracerProviderBuilder tracerProviderBuilder, ContextExtractor contextExtractor, boolean enableMdc) { + this.idGenerator = new DeterministicIdGenerator(); + this.tracerProvider = tracerProviderBuilder.setIdGenerator(idGenerator).build(); this.tracer = tracerProvider.get(INSTRUMENTATION_NAME); this.contextExtractor = contextExtractor; - this.samplingRate = samplingRate; this.enableMdc = enableMdc; } - /** - * Creates an OTel plugin using a pre-configured {@link SdkTracerProvider}. - * - * @param sdkTracerProvider the SDK tracer provider - * @param idGenerator the deterministic ID generator - */ - public OpenTelemetryDurablePlugin(SdkTracerProvider sdkTracerProvider, DeterministicIdGenerator idGenerator) { - this((TracerProvider) sdkTracerProvider, idGenerator); - } - // ─── Invocation hooks ──────────────────────────────────────────────── @Override @@ -140,10 +119,6 @@ public void onInvocationStart(InvocationInfo info) { this.executionArn = info.executionArn(); idGenerator.setExecutionArn(info.executionArn()); - // Determine sampling (consistent across all invocations of same execution) - this.sampled = SamplingUtil.shouldSampleExecution(info.executionArn(), samplingRate); - if (!sampled) return; - // Extract parent context from Lambda environment (X-Ray, W3C, etc.) var extractedParentContext = contextExtractor.extract(); @@ -162,7 +137,7 @@ public void onInvocationStart(InvocationInfo info) { @Override public void onInvocationEnd(InvocationEndInfo info) { - if (!sampled || invocationSpan == null) return; + if (invocationSpan == null) return; // End any operation spans that are still open (operations that didn't complete in this invocation) for (var entry : operationSpans.entrySet()) { @@ -196,11 +171,9 @@ public void onInvocationEnd(InvocationEndInfo info) { invocationSpan = null; // Flush spans before Lambda freezes - if (tracerProvider instanceof SdkTracerProvider sdkProvider) { - var flushResult = sdkProvider.forceFlush().join(5, java.util.concurrent.TimeUnit.SECONDS); - if (!flushResult.isSuccess()) { - logger.warn("OTel span flush failed or timed out — some spans may be lost"); - } + var flushResult = tracerProvider.forceFlush().join(5, java.util.concurrent.TimeUnit.SECONDS); + if (!flushResult.isSuccess()) { + logger.warn("OTel span flush failed or timed out — some spans may be lost"); } } @@ -208,7 +181,7 @@ public void onInvocationEnd(InvocationEndInfo info) { @Override public void onOperationStart(OperationInfo info) { - if (!sampled || info.id() == null) return; + if (info.id() == null) return; idGenerator.setNextSpanOperationId(info.id()); @@ -239,7 +212,7 @@ public void onOperationStart(OperationInfo info) { @Override public void onOperationEnd(OperationEndInfo info) { - if (!sampled || info.id() == null) return; + if (info.id() == null) return; // End the operation span that was started in onOperationStart var span = operationSpans.remove(info.id()); @@ -257,7 +230,6 @@ public void onOperationEnd(OperationEndInfo info) { @Override public void onUserFunctionStart(UserFunctionStartInfo info) { - if (!sampled) return; var key = attemptKey(info.id(), info.attempt()); // Use the operation span as parent for the attempt span @@ -295,7 +267,6 @@ public void onUserFunctionStart(UserFunctionStartInfo info) { @Override public void onUserFunctionEnd(UserFunctionEndInfo info) { - if (!sampled) return; var key = attemptKey(info.id(), info.attempt()); // Close scope first (must happen on same thread as makeCurrent) diff --git a/otel-plugin/src/main/java/software/amazon/lambda/durable/otel/SamplingUtil.java b/otel-plugin/src/main/java/software/amazon/lambda/durable/otel/SamplingUtil.java deleted file mode 100644 index 65d3e8b0b..000000000 --- a/otel-plugin/src/main/java/software/amazon/lambda/durable/otel/SamplingUtil.java +++ /dev/null @@ -1,55 +0,0 @@ -// Copyright Amazon.com, Inc. or its affiliates. All Rights Reserved. -// SPDX-License-Identifier: Apache-2.0 -package software.amazon.lambda.durable.otel; - -import java.nio.charset.StandardCharsets; - -/** - * Deterministic sampling utility for durable executions. - * - *

Uses FNV-1a hash of the execution ARN to decide whether to sample. This ensures consistent sampling across all - * invocations of the same execution — if you sample the first invocation, you sample all subsequent ones. - * - * @deprecated This is a preview API that is experimental and may be changed or removed in future releases. - */ -@Deprecated -public final class SamplingUtil { - - private SamplingUtil() {} - - // FNV-1a 32-bit constants - private static final int FNV_OFFSET_BASIS = 0x811c9dc5; - private static final int FNV_PRIME = 0x01000193; - - /** - * Determines whether an execution should be sampled based on its ARN. - * - *

Uses FNV-1a hash to distribute executions uniformly. The same ARN always produces the same result for a given - * sampling rate. - * - * @param executionArn the durable execution ARN - * @param samplingRate value between 0.0 (sample nothing) and 1.0 (sample everything) - * @return true if this execution should be sampled - */ - public static boolean shouldSampleExecution(String executionArn, double samplingRate) { - if (samplingRate >= 1.0) return true; - if (samplingRate <= 0.0) return false; - if (executionArn == null || executionArn.isEmpty()) return false; - - var hash = fnv1a32(executionArn); - // Convert to a value between 0.0 and 1.0 - var normalized = (hash & 0xFFFFFFFFL) / (double) 0x100000000L; - return normalized < samplingRate; - } - - /** FNV-1a 32-bit hash function. */ - static int fnv1a32(String input) { - var bytes = input.getBytes(StandardCharsets.UTF_8); - int hash = FNV_OFFSET_BASIS; - for (byte b : bytes) { - hash ^= (b & 0xFF); - hash *= FNV_PRIME; - } - return hash; - } -} diff --git a/otel-plugin/src/test/java/software/amazon/lambda/durable/otel/MdcSpanEnricherTest.java b/otel-plugin/src/test/java/software/amazon/lambda/durable/otel/MdcSpanEnricherTest.java index 164104a89..b601d9853 100644 --- a/otel-plugin/src/test/java/software/amazon/lambda/durable/otel/MdcSpanEnricherTest.java +++ b/otel-plugin/src/test/java/software/amazon/lambda/durable/otel/MdcSpanEnricherTest.java @@ -54,14 +54,11 @@ void inject_withNullArn_doesNotSetArn() { void plugin_withMdcEnabled_setsArnInMdc() { // Test MDC through the full plugin lifecycle (where makeCurrent is called on same thread) var spanExporter = InMemorySpanExporter.create(); - var idGenerator = new DeterministicIdGenerator(); - var tracerProvider = SdkTracerProvider.builder() - .setIdGenerator(idGenerator) - .addSpanProcessor(SimpleSpanProcessor.create(spanExporter)) - .build(); var plugin = new OpenTelemetryDurablePlugin( - tracerProvider, idGenerator, () -> io.opentelemetry.context.Context.root(), 1.0, true); + SdkTracerProvider.builder().addSpanProcessor(SimpleSpanProcessor.create(spanExporter)), + () -> io.opentelemetry.context.Context.root(), + true); plugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec-mdc-test", true)); diff --git a/otel-plugin/src/test/java/software/amazon/lambda/durable/otel/OpenTelemetryDurablePluginTest.java b/otel-plugin/src/test/java/software/amazon/lambda/durable/otel/OpenTelemetryDurablePluginTest.java index 30bd29aa5..ea8b81ed0 100644 --- a/otel-plugin/src/test/java/software/amazon/lambda/durable/otel/OpenTelemetryDurablePluginTest.java +++ b/otel-plugin/src/test/java/software/amazon/lambda/durable/otel/OpenTelemetryDurablePluginTest.java @@ -17,20 +17,15 @@ class OpenTelemetryDurablePluginTest { private InMemorySpanExporter spanExporter; private OpenTelemetryDurablePlugin plugin; - private DeterministicIdGenerator idGenerator; @BeforeEach void setUp() { spanExporter = InMemorySpanExporter.create(); - idGenerator = new DeterministicIdGenerator(); - - var tracerProvider = SdkTracerProvider.builder() - .setIdGenerator(idGenerator) - .addSpanProcessor(SimpleSpanProcessor.create(spanExporter)) - .build(); plugin = new OpenTelemetryDurablePlugin( - tracerProvider, idGenerator, () -> io.opentelemetry.context.Context.root(), 1.0, false); + SdkTracerProvider.builder().addSpanProcessor(SimpleSpanProcessor.create(spanExporter)), + () -> io.opentelemetry.context.Context.root(), + false); } @Test @@ -215,16 +210,13 @@ void operationNotCompleted_spanEndedAtInvocationEnd() { } @Test - void sampling_disabledExecution_producesNoSpans() { + void sampling_disabled_producesNoSpans() { spanExporter = InMemorySpanExporter.create(); var sampledPlugin = new OpenTelemetryDurablePlugin( SdkTracerProvider.builder() - .setIdGenerator(idGenerator) - .addSpanProcessor(SimpleSpanProcessor.create(spanExporter)) - .build(), - idGenerator, + .setSampler(io.opentelemetry.sdk.trace.samplers.Sampler.alwaysOff()) + .addSpanProcessor(SimpleSpanProcessor.create(spanExporter)), () -> io.opentelemetry.context.Context.root(), - 0.0, // 0% sampling — nothing should be traced false); sampledPlugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec1", true)); diff --git a/otel-plugin/src/test/java/software/amazon/lambda/durable/otel/SamplingUtilTest.java b/otel-plugin/src/test/java/software/amazon/lambda/durable/otel/SamplingUtilTest.java deleted file mode 100644 index ed7817321..000000000 --- a/otel-plugin/src/test/java/software/amazon/lambda/durable/otel/SamplingUtilTest.java +++ /dev/null @@ -1,74 +0,0 @@ -// Copyright Amazon.com, Inc. or its affiliates. All Rights Reserved. -// SPDX-License-Identifier: Apache-2.0 -package software.amazon.lambda.durable.otel; - -import static org.junit.jupiter.api.Assertions.*; - -import org.junit.jupiter.api.Test; - -class SamplingUtilTest { - - @Test - void samplingRate1_alwaysSamples() { - assertTrue(SamplingUtil.shouldSampleExecution("arn:exec1", 1.0)); - assertTrue(SamplingUtil.shouldSampleExecution("arn:exec2", 1.0)); - assertTrue(SamplingUtil.shouldSampleExecution("anything", 1.0)); - } - - @Test - void samplingRate0_neverSamples() { - assertFalse(SamplingUtil.shouldSampleExecution("arn:exec1", 0.0)); - assertFalse(SamplingUtil.shouldSampleExecution("arn:exec2", 0.0)); - } - - @Test - void deterministic_sameArnSameResult() { - var arn = "arn:aws:lambda:us-east-1:123:function:test:$LATEST/durable/exec1"; - var result1 = SamplingUtil.shouldSampleExecution(arn, 0.5); - var result2 = SamplingUtil.shouldSampleExecution(arn, 0.5); - - assertEquals(result1, result2, "Same ARN should always produce same sampling decision"); - } - - @Test - void distribution_approximatesRate() { - // With enough samples, the fraction sampled should approximate the rate - int sampled = 0; - int total = 10000; - for (int i = 0; i < total; i++) { - if (SamplingUtil.shouldSampleExecution("arn:exec-" + i, 0.5)) { - sampled++; - } - } - - double actualRate = (double) sampled / total; - // Allow 5% tolerance - assertTrue(actualRate > 0.45 && actualRate < 0.55, "Sampling rate should be ~50%, got " + actualRate); - } - - @Test - void nullArn_returnsFalse() { - assertFalse(SamplingUtil.shouldSampleExecution(null, 0.5)); - } - - @Test - void emptyArn_returnsFalse() { - assertFalse(SamplingUtil.shouldSampleExecution("", 0.5)); - } - - @Test - void fnv1a32_isDeterministic() { - var hash1 = SamplingUtil.fnv1a32("test-input"); - var hash2 = SamplingUtil.fnv1a32("test-input"); - - assertEquals(hash1, hash2); - } - - @Test - void fnv1a32_differentInputsDifferentHashes() { - var hash1 = SamplingUtil.fnv1a32("input-a"); - var hash2 = SamplingUtil.fnv1a32("input-b"); - - assertNotEquals(hash1, hash2); - } -} diff --git a/sdk-integration-tests/src/test/java/software/amazon/lambda/durable/OtelPluginIntegrationTest.java b/sdk-integration-tests/src/test/java/software/amazon/lambda/durable/OtelPluginIntegrationTest.java index 66b0c74f1..030b069b3 100644 --- a/sdk-integration-tests/src/test/java/software/amazon/lambda/durable/OtelPluginIntegrationTest.java +++ b/sdk-integration-tests/src/test/java/software/amazon/lambda/durable/OtelPluginIntegrationTest.java @@ -15,7 +15,6 @@ import org.junit.jupiter.api.Test; import software.amazon.lambda.durable.config.StepConfig; import software.amazon.lambda.durable.model.ExecutionStatus; -import software.amazon.lambda.durable.otel.DeterministicIdGenerator; import software.amazon.lambda.durable.otel.OpenTelemetryDurablePlugin; import software.amazon.lambda.durable.retry.RetryStrategies; import software.amazon.lambda.durable.testing.LocalDurableTestRunner; @@ -27,21 +26,16 @@ class OtelPluginIntegrationTest { private InMemorySpanExporter spanExporter; - private DeterministicIdGenerator idGenerator; private DurableConfig otelConfig; @BeforeEach void setUp() { spanExporter = InMemorySpanExporter.create(); - idGenerator = new DeterministicIdGenerator(); - - var tracerProvider = SdkTracerProvider.builder() - .setIdGenerator(idGenerator) - .addSpanProcessor(SimpleSpanProcessor.create(spanExporter)) - .build(); var plugin = new OpenTelemetryDurablePlugin( - tracerProvider, idGenerator, () -> io.opentelemetry.context.Context.root(), 1.0, false); + SdkTracerProvider.builder().addSpanProcessor(SimpleSpanProcessor.create(spanExporter)), + () -> io.opentelemetry.context.Context.root(), + false); otelConfig = DurableConfig.builder().withPlugins(plugin).build(); } @@ -218,16 +212,15 @@ void failedStep_producesErrorSpan() { } @Test - void sampling_zeroRate_producesNoSpans() { + void sampling_off_producesNoSpans() { var sampledExporter = InMemorySpanExporter.create(); - var sampledIdGen = new DeterministicIdGenerator(); - var sampledProvider = SdkTracerProvider.builder() - .setIdGenerator(sampledIdGen) - .addSpanProcessor(SimpleSpanProcessor.create(sampledExporter)) - .build(); var noSamplePlugin = new OpenTelemetryDurablePlugin( - sampledProvider, sampledIdGen, () -> io.opentelemetry.context.Context.root(), 0.0, false); + SdkTracerProvider.builder() + .setSampler(io.opentelemetry.sdk.trace.samplers.Sampler.alwaysOff()) + .addSpanProcessor(SimpleSpanProcessor.create(sampledExporter)), + () -> io.opentelemetry.context.Context.root(), + false); var noSampleConfig = DurableConfig.builder().withPlugins(noSamplePlugin).build();