From 562f9c33cdffdfcda79f836fc57f29b4b4009c1e Mon Sep 17 00:00:00 2001 From: Frank Chen Date: Fri, 2 Oct 2026 22:54:30 +0000 Subject: [PATCH 1/7] feat(plugin): prepare chained-invoke propagation metadata --- docs/advanced/propagation-metadata.md | 39 ++++ otel-plugin/README.md | 6 + .../durable/otel/ExecutionOtelPlugin.java | 20 ++ .../durable/otel/InvocationOtelPlugin.java | 20 ++ .../durable/otel/OtelPropagationMetadata.java | 18 ++ .../otel/OtelPropagationMetadataTest.java | 196 ++++++++++++++++++ .../PropagationCollectorDiagnosticsTest.java | 60 ++++++ .../plugin/DurableExecutionPlugin.java | 21 ++ .../lambda/durable/plugin/PluginRunner.java | 52 +++++ .../durable/plugin/PropagationInput.java | 45 ++++ .../durable/plugin/PropagationMetadata.java | 19 ++ .../PropagationMetadataCollectorTest.java | 111 ++++++++++ 12 files changed, 607 insertions(+) create mode 100644 docs/advanced/propagation-metadata.md create mode 100644 otel-plugin/src/main/java/software/amazon/lambda/durable/otel/OtelPropagationMetadata.java create mode 100644 otel-plugin/src/test/java/software/amazon/lambda/durable/otel/OtelPropagationMetadataTest.java create mode 100644 otel-plugin/src/test/java/software/amazon/lambda/durable/otel/PropagationCollectorDiagnosticsTest.java create mode 100644 sdk/src/main/java/software/amazon/lambda/durable/plugin/PropagationInput.java create mode 100644 sdk/src/main/java/software/amazon/lambda/durable/plugin/PropagationMetadata.java create mode 100644 sdk/src/test/java/software/amazon/lambda/durable/plugin/PropagationMetadataCollectorTest.java diff --git a/docs/advanced/propagation-metadata.md b/docs/advanced/propagation-metadata.md new file mode 100644 index 000000000..98db16349 --- /dev/null +++ b/docs/advanced/propagation-metadata.md @@ -0,0 +1,39 @@ +# Chained-invoke propagation groundwork + +This draft prepares a plugin contract and pure OTel producers. It does **not** send propagation metadata or complete +chained-invoke propagation. Current public Lambda client models do not expose `ChainedInvokeOptions.XAmznTraceId` +or `DistributedMapOptions`. Production invoke START requests do not call this hook, and no serializer or raw HTTP +workaround is included. + +`PropagationInput` is an immutable SDK-owned snapshot with `executionArn`, `operationId`, optional `parentOperationId`, +and `targetFunctionName`. `PropagationMetadata` is an immutable SDK-owned contribution with optional `xAmznTraceId`. +Both types use only Java strings; generated Lambda and OTel model types do not enter the plugin interface. +The typed default hook adapts a String-only producer overload, so bundled plugin layers do not reference newly added +core classes. Older cores continue to load and run those layers; new metadata collection requires a newer core, +without changing registration/provider API versions or rejecting previously working plugins. + +The optional synchronous `providePropagationMetadata` default method returns null. Existing plugins keep compiling +and linking. The dispatcher supplies the same immutable snapshot in configured order. The first non-null supported +member wins; equal values do not conflict. Different later values warn with the winning and conflicting plugin +identities and a conflict count, without logging the header. Null contributions abstain; blank headers and ordinary +hook failures are logged and skipped. Java's return type rules out asynchronous and foreign result objects. +The collector preserves the current event dispatch policy: `Exception` (including `CancellationException`) is +contained, while `Error` propagates. A failing logging backend does not interrupt collection. + +Both OTel views encode their resolved canonical trace ID, operation span ID and sampling decision as +`Root=1-<8 hex>-<24 hex>;Parent=<16 hex>;Sampled=0|1`. They require an active invocation with matching execution +ownership. An observed operation's context is authoritative, including an invocation-view continuation segment. +Before a new operation starts, the producer derives the same initial-operation ID the normal start hook will use. +It creates no spans, samples nothing again, and does not alter the existing provider, resource, upstream-parent or +fallback behavior. An explicit upstream `Sampled=0` still produces metadata with `Sampled=0`. + +Remaining work before this can implement the feature: + +1. Supported public client models and generated serialization for `ChainedInvokeOptions.XAmznTraceId` and + `DistributedMapOptions`. +2. Backend rollout of the reviewed flat X-Ray header contract. +3. Consuming the collector on the real operation START path, including replay and failed-checkpoint integration. +4. Deployed trace topology and propagation validation. + +Keep the implementation draft while these dependencies block completion. This is related to Java issue #764; +it does not close that issue. Runtime inbound headers and plugin-instance lifetime remain separate changes. diff --git a/otel-plugin/README.md b/otel-plugin/README.md index c51032aaa..aceded49f 100644 --- a/otel-plugin/README.md +++ b/otel-plugin/README.md @@ -40,6 +40,12 @@ If you configure your own `SdkTracerProviderBuilder`, add the OpenTelemetry SDK ``` +## Chained-invoke propagation groundwork + +The SDK exposes a draft, synchronous metadata contract and both views provide pure X-Ray metadata producers. +Production invoke START integration and supported public client models are still pending; this does not enable +outbound propagation. See [the scope and remaining dependencies](../docs/advanced/propagation-metadata.md). + ## Quick Start using X-Ray/CloudWatch Tracing (ADOT Java Agent) 1. Add the ADOT Lambda Layer to your function diff --git a/otel-plugin/src/main/java/software/amazon/lambda/durable/otel/ExecutionOtelPlugin.java b/otel-plugin/src/main/java/software/amazon/lambda/durable/otel/ExecutionOtelPlugin.java index 09214055c..c9381749a 100644 --- a/otel-plugin/src/main/java/software/amazon/lambda/durable/otel/ExecutionOtelPlugin.java +++ b/otel-plugin/src/main/java/software/amazon/lambda/durable/otel/ExecutionOtelPlugin.java @@ -346,6 +346,26 @@ public void onInvocationEnd(InvocationEndInfo info) { } } + /** Produces metadata without creating spans or changing the configured provider/resource. */ + @Override + public String providePropagationMetadata( + String executionArn, String operationId, String parentOperationId, String targetFunctionName) { + if (!tracingEnabled || !executionArn.equals(durableExecutionArn)) return null; + var trace = executionTrace; + if (trace == null) return null; + // An observed operation's actual context is authoritative, including invocation-view continuation segments. + var context = operationContexts.get(operationId); + if (context == null) { + // Before a new operation starts, use the same deterministic ID its initial span will receive. + context = SpanContext.create( + trace.traceId(), + idGenerator.generateSpanIdForOperation(durableExecutionArn, operationId), + effectiveTraceFlags(), + effectiveTraceState()); + } + return OtelPropagationMetadata.fromContext(context); + } + // ─── Operation hooks ───────────────────────────────────────────────── @Override diff --git a/otel-plugin/src/main/java/software/amazon/lambda/durable/otel/InvocationOtelPlugin.java b/otel-plugin/src/main/java/software/amazon/lambda/durable/otel/InvocationOtelPlugin.java index e89804bdc..d458e1a3d 100644 --- a/otel-plugin/src/main/java/software/amazon/lambda/durable/otel/InvocationOtelPlugin.java +++ b/otel-plugin/src/main/java/software/amazon/lambda/durable/otel/InvocationOtelPlugin.java @@ -352,6 +352,26 @@ public void onInvocationEnd(InvocationEndInfo info) { } } + /** Produces metadata without creating spans or changing the configured provider/resource. */ + @Override + public String providePropagationMetadata( + String executionArn, String operationId, String parentOperationId, String targetFunctionName) { + if (!tracingEnabled || !executionArn.equals(durableExecutionArn)) return null; + var trace = executionTrace; + if (trace == null) return null; + // An observed operation's actual context is authoritative, including invocation-view continuation segments. + var context = operationContexts.get(operationId); + if (context == null) { + // Before a new operation starts, use the same deterministic ID its initial span will receive. + context = SpanContext.create( + trace.traceId(), + idGenerator.generateSpanIdForOperation(durableExecutionArn, operationId), + effectiveTraceFlags(), + effectiveTraceState()); + } + return OtelPropagationMetadata.fromContext(context); + } + // ─── Operation hooks ───────────────────────────────────────────────── @Override diff --git a/otel-plugin/src/main/java/software/amazon/lambda/durable/otel/OtelPropagationMetadata.java b/otel-plugin/src/main/java/software/amazon/lambda/durable/otel/OtelPropagationMetadata.java new file mode 100644 index 000000000..88af59068 --- /dev/null +++ b/otel-plugin/src/main/java/software/amazon/lambda/durable/otel/OtelPropagationMetadata.java @@ -0,0 +1,18 @@ +// Copyright Amazon.com, Inc. or its affiliates. All Rights Reserved. +// SPDX-License-Identifier: Apache-2.0 +package software.amazon.lambda.durable.otel; + +import io.opentelemetry.api.trace.SpanContext; + +/** Pure encoding of an existing operation context into the reviewed X-Ray carrier shape. */ +final class OtelPropagationMetadata { + private OtelPropagationMetadata() {} + + static String fromContext(SpanContext context) { + if (context == null || !context.isValid()) return null; + var trace = context.getTraceId(); + var header = "Root=1-" + trace.substring(0, 8) + "-" + trace.substring(8) + ";Parent=" + context.getSpanId() + + ";Sampled=" + (context.isSampled() ? "1" : "0"); + return header; + } +} diff --git a/otel-plugin/src/test/java/software/amazon/lambda/durable/otel/OtelPropagationMetadataTest.java b/otel-plugin/src/test/java/software/amazon/lambda/durable/otel/OtelPropagationMetadataTest.java new file mode 100644 index 000000000..1979b1071 --- /dev/null +++ b/otel-plugin/src/test/java/software/amazon/lambda/durable/otel/OtelPropagationMetadataTest.java @@ -0,0 +1,196 @@ +// 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 io.opentelemetry.api.common.AttributeKey; +import io.opentelemetry.api.common.Attributes; +import io.opentelemetry.api.trace.Span; +import io.opentelemetry.api.trace.SpanContext; +import io.opentelemetry.api.trace.TraceFlags; +import io.opentelemetry.api.trace.TraceState; +import io.opentelemetry.sdk.resources.Resource; +import io.opentelemetry.sdk.testing.exporter.InMemorySpanExporter; +import io.opentelemetry.sdk.trace.SdkTracerProvider; +import io.opentelemetry.sdk.trace.export.SimpleSpanProcessor; +import io.opentelemetry.sdk.trace.samplers.Sampler; +import java.time.Instant; +import java.util.List; +import java.util.concurrent.atomic.AtomicBoolean; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.CsvSource; +import org.junit.jupiter.params.provider.ValueSource; +import software.amazon.lambda.durable.plugin.DurableExecutionPlugin; +import software.amazon.lambda.durable.plugin.InvocationEndInfo; +import software.amazon.lambda.durable.plugin.InvocationInfo; +import software.amazon.lambda.durable.plugin.InvocationStatus; +import software.amazon.lambda.durable.plugin.OperationEndInfo; +import software.amazon.lambda.durable.plugin.OperationInfo; +import software.amazon.lambda.durable.plugin.PluginRunner; +import software.amazon.lambda.durable.plugin.PropagationInput; + +class OtelPropagationMetadataTest { + private static final String TRACE = "6955b900123456789012345678901234"; + private static final String ARN = "arn:aws:lambda:us-east-1:123456789012:function:parent/durable-execution/test/id"; + private static final Instant START = Instant.parse("2026-10-02T00:00:00Z"); + private static final PropagationInput INPUT = + new PropagationInput(ARN, "invoke-op", "parent-context", "child:live"); + private static final AttributeKey SERVICE = AttributeKey.stringKey("service.name"); + + @ParameterizedTest + @CsvSource({"true,true", "true,false", "false,true", "false,false"}) + void collectorAndProducerMatchTheActualOperationAndSampling(boolean executionView, boolean sampled) { + var fixture = fixture(executionView, sampled, new AtomicBoolean()); + var runner = fixture.runner(); + assertNull(runner.providePropagationMetadata(INPUT)); + var ambient = Span.wrap(SpanContext.create( + "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa", + "bbbbbbbbbbbbbbbb", + TraceFlags.getSampled(), + TraceState.getDefault())); + try (var ignored = ambient.makeCurrent()) { + start(runner, ARN, true); + var planned = runner.providePropagationMetadata(INPUT).xAmznTraceId(); + assertTrue(fixture.exporter().getFinishedSpanItems().isEmpty(), "Metadata collection creates no span"); + runner.onOperationStart(operation(false)); + assertEquals(planned, runner.providePropagationMetadata(INPUT).xAmznTraceId()); + runner.onOperationEnd(operationEnd()); + end(runner, ARN, true, InvocationStatus.SUCCEEDED); + assertHeader( + planned, + sampled, + new DeterministicIdGenerator().generateSpanIdForOperation(ARN, INPUT.operationId())); + assertExportedOperation(fixture, planned, sampled); + } + assertNull(runner.providePropagationMetadata(INPUT), "Inactive invocation has no contribution"); + } + + @ParameterizedTest + @ValueSource(booleans = {true, false}) + void observedContinuationUsesItsActualSpanContext(boolean executionView) { + var fixture = fixture(executionView, true, new AtomicBoolean()); + start(fixture.runner(), ARN, false); + fixture.runner().onOperationStart(operation(true)); + var header = fixture.runner().providePropagationMetadata(INPUT).xAmznTraceId(); + fixture.runner().onOperationEnd(operationEnd()); + end(fixture.runner(), ARN, false, InvocationStatus.SUCCEEDED); + assertExportedOperation(fixture, header, true); + } + + @ParameterizedTest + @ValueSource(booleans = {true, false}) + void resumeIdentityAndExecutionOwnershipRemainStable(boolean executionView) { + var first = fixture(executionView, true, new AtomicBoolean()); + start(first.runner(), ARN, true); + var initial = first.runner().providePropagationMetadata(INPUT).xAmznTraceId(); + end(first.runner(), ARN, true, InvocationStatus.PENDING); + var resumed = fixture(executionView, true, new AtomicBoolean()); + start(resumed.runner(), ARN, false); + assertEquals(initial, resumed.runner().providePropagationMetadata(INPUT).xAmznTraceId()); + assertNull( + resumed.runner().providePropagationMetadata(new PropagationInput("other-arn", "op", null, "target"))); + end(resumed.runner(), ARN, false, InvocationStatus.SUCCEEDED); + start(resumed.runner(), "other-arn", true); + assertNull(resumed.runner().providePropagationMetadata(INPUT)); + end(resumed.runner(), "other-arn", true, InvocationStatus.SUCCEEDED); + } + + @ParameterizedTest + @ValueSource(booleans = {true, false}) + void failedSetupCannotExposePreviousExecutionContext(boolean executionView) { + var fail = new AtomicBoolean(); + var fixture = fixture(executionView, true, fail); + start(fixture.runner(), ARN, true); + assertNotNull(fixture.runner().providePropagationMetadata(INPUT)); + end(fixture.runner(), ARN, true, InvocationStatus.SUCCEEDED); + fail.set(true); + start(fixture.runner(), "other-arn", true); // The runner contains the extractor's ordinary failure. + assertNull(fixture.runner().providePropagationMetadata(INPUT)); + assertNull( + fixture.runner().providePropagationMetadata(new PropagationInput("other-arn", "op", null, "target"))); + end(fixture.runner(), "other-arn", true, InvocationStatus.FAILED); + } + + private static void assertHeader(String header, boolean sampled, String parentId) { + assertEquals( + "Root=1-6955b900-123456789012345678901234;Parent=" + parentId + ";Sampled=" + (sampled ? "1" : "0"), + header); + } + + private static void assertExportedOperation(Fixture fixture, String header, boolean sampled) { + var spans = fixture.exporter().getFinishedSpanItems(); + if (!sampled) { + assertTrue(spans.isEmpty()); + return; + } + assertEquals(3, spans.size(), "Only the normal operation, Invocation and Workflow spans are exported"); + var operation = spans.stream() + .filter(span -> span.getName().equals("invoke")) + .findFirst() + .orElseThrow(); + assertHeader(header, true, operation.getSpanId()); + assertEquals(TRACE, operation.getTraceId()); + assertTrue(spans.stream() + .allMatch(span -> "configured-service".equals(span.getResource().getAttribute(SERVICE)))); + } + + private static Fixture fixture(boolean executionView, boolean sampled, AtomicBoolean fail) { + var exporter = InMemorySpanExporter.create(); + var builder = SdkTracerProvider.builder() + .setResource(Resource.create(Attributes.of(SERVICE, "configured-service"))) + .setSampler(sampled ? Sampler.alwaysOff() : Sampler.alwaysOn()) + .addSpanProcessor(SimpleSpanProcessor.create(exporter)); + var config = OtelPluginConfig.builder() + .enableMdc(false) + .contextExtractor(() -> { + if (fail.get()) throw new IllegalStateException("extractor failed"); + return new ExtractedContext( + TRACE, + "1234567890123456", + sampled ? ExtractedContext.Sampling.SAMPLED : ExtractedContext.Sampling.NOT_SAMPLED); + }) + .build(); + DurableExecutionPlugin plugin = + executionView ? new ExecutionOtelPlugin(builder, config) : new InvocationOtelPlugin(builder, config); + return new Fixture(new PluginRunner(List.of(new DurableExecutionPlugin() {}, plugin)), exporter); + } + + private static OperationInfo operation(boolean replay) { + return new OperationInfo( + INPUT.operationId(), + "invoke", + "CHAINED_INVOKE", + null, + INPUT.parentOperationId(), + START, + null, + "STARTED", + replay); + } + + private static OperationEndInfo operationEnd() { + return new OperationEndInfo( + INPUT.operationId(), + "invoke", + "CHAINED_INVOKE", + null, + INPUT.parentOperationId(), + START, + START.plusSeconds(1), + "SUCCEEDED", + null, + false, + null); + } + + private static void start(PluginRunner runner, String arn, boolean first) { + runner.onInvocationStart(new InvocationInfo("request", arn, first, START)); + } + + private static void end(PluginRunner runner, String arn, boolean first, InvocationStatus status) { + runner.onInvocationEnd(new InvocationEndInfo("request", arn, first, status, null)); + } + + private record Fixture(PluginRunner runner, InMemorySpanExporter exporter) {} +} diff --git a/otel-plugin/src/test/java/software/amazon/lambda/durable/otel/PropagationCollectorDiagnosticsTest.java b/otel-plugin/src/test/java/software/amazon/lambda/durable/otel/PropagationCollectorDiagnosticsTest.java new file mode 100644 index 000000000..e18794400 --- /dev/null +++ b/otel-plugin/src/test/java/software/amazon/lambda/durable/otel/PropagationCollectorDiagnosticsTest.java @@ -0,0 +1,60 @@ +// 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 ch.qos.logback.classic.Logger; +import ch.qos.logback.classic.spi.ILoggingEvent; +import ch.qos.logback.core.read.ListAppender; +import java.util.List; +import org.junit.jupiter.api.Test; +import org.slf4j.LoggerFactory; +import software.amazon.lambda.durable.plugin.DurableExecutionPlugin; +import software.amazon.lambda.durable.plugin.PluginRunner; +import software.amazon.lambda.durable.plugin.PropagationInput; +import software.amazon.lambda.durable.plugin.PropagationMetadata; + +class PropagationCollectorDiagnosticsTest { + @Test + void differentValuesNameBothOwnersAndCountConflictsWithoutLoggingHeaders() { + var logger = (Logger) LoggerFactory.getLogger(PluginRunner.class); + var appender = new ListAppender(); + appender.start(); + logger.addAppender(appender); + try { + var runner = new PluginRunner(List.of( + plugin("private-first"), + plugin("private-first"), + plugin("private-second"), + plugin("private-third"))); + assertEquals( + "private-first", + runner.providePropagationMetadata(new PropagationInput("arn", "op", null, "target")) + .xAmznTraceId()); + var warnings = appender.list.stream() + .map(ILoggingEvent::getFormattedMessage) + .toList(); + assertEquals(2, warnings.size(), "Equal values are not conflicts"); + assertTrue(warnings.get(0).contains("[2]")); + assertTrue(warnings.get(0).contains("[0]")); + assertTrue(warnings.get(0).contains("conflict count: 1")); + assertTrue(warnings.get(1).contains("[3]")); + assertTrue(warnings.get(1).contains("[0]")); + assertTrue(warnings.get(1).contains("conflict count: 2")); + assertTrue(warnings.stream().noneMatch(message -> message.contains("private-"))); + } finally { + logger.detachAppender(appender); + appender.stop(); + } + } + + private static DurableExecutionPlugin plugin(String header) { + return new DurableExecutionPlugin() { + @Override + public PropagationMetadata providePropagationMetadata(PropagationInput input) { + return new PropagationMetadata(header); + } + }; + } +} diff --git a/sdk/src/main/java/software/amazon/lambda/durable/plugin/DurableExecutionPlugin.java b/sdk/src/main/java/software/amazon/lambda/durable/plugin/DurableExecutionPlugin.java index e2f8a46df..fa5603c5c 100644 --- a/sdk/src/main/java/software/amazon/lambda/durable/plugin/DurableExecutionPlugin.java +++ b/sdk/src/main/java/software/amazon/lambda/durable/plugin/DurableExecutionPlugin.java @@ -15,6 +15,27 @@ */ public interface DurableExecutionPlugin { + /** + * Optionally supplies operation-level trace propagation metadata synchronously. Return null to abstain. This + * groundwork hook is not yet called by production invoke START requests; client-model support and backend rollout + * are pending. Implementations must not mutate execution state or create side effects to produce metadata. + */ + default PropagationMetadata providePropagationMetadata(PropagationInput input) { + var header = providePropagationMetadata( + input.executionArn(), input.operationId(), input.parentOperationId(), input.targetFunctionName()); + return header != null ? new PropagationMetadata(header) : null; + } + + /** + * Optional layer-compatible X-Ray producer capability. String-only linkage lets a new plugin layer keep loading on + * older cores that do not know PropagationInput/PropagationMetadata; those cores simply never call this hook. The + * typed hook above is the preferred core contract and adapts this capability into its immutable result. + */ + default String providePropagationMetadata( + String executionArn, String operationId, String parentOperationId, String targetFunctionName) { + return null; + } + // ─── Invocation-level hooks ────────────────────────────────────────── /** diff --git a/sdk/src/main/java/software/amazon/lambda/durable/plugin/PluginRunner.java b/sdk/src/main/java/software/amazon/lambda/durable/plugin/PluginRunner.java index e3a5707c4..d5f18d8b2 100644 --- a/sdk/src/main/java/software/amazon/lambda/durable/plugin/PluginRunner.java +++ b/sdk/src/main/java/software/amazon/lambda/durable/plugin/PluginRunner.java @@ -2,6 +2,8 @@ // SPDX-License-Identifier: Apache-2.0 package software.amazon.lambda.durable.plugin; +import static java.util.Objects.requireNonNull; + import java.util.Collections; import java.util.List; import java.util.function.Consumer; @@ -42,6 +44,56 @@ public List getPlugins() { return plugins; } + /** + * Collects supported metadata in configured order on the caller's thread. First non-null member wins; matching + * later values are harmless. Ordinary plugin/invalid-result failures are logged and skipped, matching the event + * hooks' Exception containment policy; Errors continue to propagate. No production START path calls this yet. + */ + public PropagationMetadata providePropagationMetadata(PropagationInput input) { + requireNonNull(input, "input"); + String selected = null; + String owner = null; + var conflicts = 0; + for (var index = 0; index < plugins.size(); index++) { + var plugin = plugins.get(index); + var identity = plugin.getClass().getName() + "[" + index + "]"; + try { + var candidate = propagationHeader(plugin, input); + if (candidate == null) continue; + if (selected == null) { + selected = candidate; + owner = identity; + } else if (!selected.equals(candidate)) { + warnPropagation( + "Conflicting xAmznTraceId from plugin {}; retaining plugin {} (conflict count: {})", + identity, + owner, + ++conflicts); + } + } catch (Exception failure) { + warnPropagation("Propagation metadata from plugin {} failed; skipping contribution", identity, failure); + } + } + return selected != null ? new PropagationMetadata(selected) : null; + } + + private static String propagationHeader(DurableExecutionPlugin plugin, PropagationInput input) { + var metadata = plugin.providePropagationMetadata(input); + var header = metadata != null ? metadata.xAmznTraceId() : null; + if (header != null && header.isBlank()) { + throw new IllegalArgumentException("xAmznTraceId must not be blank"); + } + return header; + } + + private static void warnPropagation(String message, Object... arguments) { + try { + logger.warn(message, arguments); + } catch (RuntimeException ignored) { + // A failed logging backend cannot change optional metadata collection or the handler's result. + } + } + // ─── Event hooks ───────────────────────────────────────────────────── /** Calls a void hook on all plugins, swallowing any errors. */ diff --git a/sdk/src/main/java/software/amazon/lambda/durable/plugin/PropagationInput.java b/sdk/src/main/java/software/amazon/lambda/durable/plugin/PropagationInput.java new file mode 100644 index 000000000..e3dacfc46 --- /dev/null +++ b/sdk/src/main/java/software/amazon/lambda/durable/plugin/PropagationInput.java @@ -0,0 +1,45 @@ +// Copyright Amazon.com, Inc. or its affiliates. All Rights Reserved. +// SPDX-License-Identifier: Apache-2.0 +package software.amazon.lambda.durable.plugin; + +import static java.util.Objects.requireNonNull; + +/** Immutable SDK-owned description of an operation requesting chained-invoke propagation metadata. */ +public final class PropagationInput { + private final String executionArn; + private final String operationId; + private final String parentOperationId; + private final String targetFunctionName; + + public PropagationInput( + String executionArn, String operationId, String parentOperationId, String targetFunctionName) { + this.executionArn = required(executionArn, "executionArn"); + this.operationId = required(operationId, "operationId"); + this.parentOperationId = parentOperationId; + this.targetFunctionName = required(targetFunctionName, "targetFunctionName"); + } + + public String executionArn() { + return executionArn; + } + + public String operationId() { + return operationId; + } + + public String parentOperationId() { + return parentOperationId; + } + + public String targetFunctionName() { + return targetFunctionName; + } + + private static String required(String value, String name) { + requireNonNull(value, name); + if (value.isBlank()) { + throw new IllegalArgumentException(name + " must not be blank"); + } + return value; + } +} diff --git a/sdk/src/main/java/software/amazon/lambda/durable/plugin/PropagationMetadata.java b/sdk/src/main/java/software/amazon/lambda/durable/plugin/PropagationMetadata.java new file mode 100644 index 000000000..1b5fededd --- /dev/null +++ b/sdk/src/main/java/software/amazon/lambda/durable/plugin/PropagationMetadata.java @@ -0,0 +1,19 @@ +// Copyright Amazon.com, Inc. or its affiliates. All Rights Reserved. +// SPDX-License-Identifier: Apache-2.0 +package software.amazon.lambda.durable.plugin; + +/** + * Immutable SDK-owned propagation contribution. This is not a generated Lambda request model and is not currently + * serialized on the production invoke START path. A null header means the plugin has no contribution. + */ +public final class PropagationMetadata { + private final String xAmznTraceId; + + public PropagationMetadata(String xAmznTraceId) { + this.xAmznTraceId = xAmznTraceId; + } + + public String xAmznTraceId() { + return xAmznTraceId; + } +} diff --git a/sdk/src/test/java/software/amazon/lambda/durable/plugin/PropagationMetadataCollectorTest.java b/sdk/src/test/java/software/amazon/lambda/durable/plugin/PropagationMetadataCollectorTest.java new file mode 100644 index 000000000..9f6a7a7ab --- /dev/null +++ b/sdk/src/test/java/software/amazon/lambda/durable/plugin/PropagationMetadataCollectorTest.java @@ -0,0 +1,111 @@ +// Copyright Amazon.com, Inc. or its affiliates. All Rights Reserved. +// SPDX-License-Identifier: Apache-2.0 +package software.amazon.lambda.durable.plugin; + +import static org.junit.jupiter.api.Assertions.*; + +import java.time.Instant; +import java.util.ArrayList; +import java.util.List; +import java.util.concurrent.CancellationException; +import java.util.concurrent.atomic.AtomicInteger; +import org.junit.jupiter.api.Test; + +class PropagationMetadataCollectorTest { + private static final PropagationInput INPUT = new PropagationInput("arn", "operation", "parent", "target:alias"); + + @Test + void legacyPluginsRemainNoOpForMetadataAndKeepTheirHooks() { + var calls = new AtomicInteger(); + var legacy = new DurableExecutionPlugin() { + @Override + public void onInvocationStart(InvocationInfo info) { + calls.incrementAndGet(); + } + }; + var runner = new PluginRunner(List.of(legacy)); + assertNull(runner.providePropagationMetadata(INPUT)); + runner.onInvocationStart(new InvocationInfo("request", "arn", true, Instant.EPOCH)); + assertEquals(1, calls.get()); + assertNull(PluginRunner.noOp().providePropagationMetadata(INPUT)); + } + + @Test + void configuredOrderFirstNonNullWinsWithoutSkippingOtherPlugins() { + var seen = new ArrayList(); + var thread = Thread.currentThread(); + var first = new DurableExecutionPlugin() { + @Override + public PropagationMetadata providePropagationMetadata(PropagationInput input) { + assertSame(INPUT, input); + assertSame(thread, Thread.currentThread()); + seen.add(input.executionArn() + ":" + input.operationId() + ":" + input.parentOperationId() + ":" + + input.targetFunctionName()); + return new PropagationMetadata("first"); + } + }; + var runner = new PluginRunner(List.of( + new DurableExecutionPlugin() {}, + first, + plugin("equal", "first", seen), + plugin("conflict", "later", seen), + plugin("last", null, seen))); + assertEquals("first", runner.providePropagationMetadata(INPUT).xAmznTraceId()); + assertEquals(List.of("arn:operation:parent:target:alias", "equal", "conflict", "last"), seen); + } + + @Test + void ordinaryFailuresInvalidContributionsAndCancellationKeepHealthyPluginsRunning() { + var seen = new ArrayList(); + var failure = new DurableExecutionPlugin() { + @Override + public PropagationMetadata providePropagationMetadata(PropagationInput input) { + throw new IllegalStateException("optional instrumentation failed"); + } + }; + var cancellation = new DurableExecutionPlugin() { + @Override + public PropagationMetadata providePropagationMetadata(PropagationInput input) { + throw new CancellationException("optional instrumentation cancelled"); + } + }; + var runner = new PluginRunner(List.of( + failure, cancellation, plugin("invalid", " ", seen), plugin("healthy", "healthy-header", seen))); + assertEquals("healthy-header", runner.providePropagationMetadata(INPUT).xAmznTraceId()); + assertEquals(List.of("invalid", "healthy"), seen); + } + + @Test + void errorPropagationMatchesExistingPluginBoundary() { + var fatal = new AssertionError("not an ordinary Exception"); + var plugin = new DurableExecutionPlugin() { + @Override + public PropagationMetadata providePropagationMetadata(PropagationInput input) { + throw fatal; + } + }; + assertSame( + fatal, + assertThrows( + AssertionError.class, + () -> new PluginRunner(List.of(plugin)).providePropagationMetadata(INPUT))); + } + + @Test + void invalidSdkInputCannotDescribeAnOperation() { + assertThrows(NullPointerException.class, () -> new PropagationInput(null, "op", null, "target")); + assertThrows(IllegalArgumentException.class, () -> new PropagationInput("arn", "", null, "target")); + assertThrows(IllegalArgumentException.class, () -> new PropagationInput("arn", "op", null, " ")); + } + + private static DurableExecutionPlugin plugin(String identity, String header, List calls) { + return new DurableExecutionPlugin() { + @Override + public PropagationMetadata providePropagationMetadata(PropagationInput input) { + assertSame(INPUT, input); + calls.add(identity); + return new PropagationMetadata(header); + } + }; + } +} From 98ad579eda9426bc1a11c365b15aee35ec68db01 Mon Sep 17 00:00:00 2001 From: Frank Chen <65260095+zhongkechen@users.noreply.github.com> Date: Mon, 5 Oct 2026 12:07:31 -0700 Subject: [PATCH 2/7] feat: propagate operation trace context on chained invoke START --- docs/advanced/propagation-metadata.md | 81 ++-- otel-plugin/README.md | 9 +- .../InvokePropagationIntegrationTest.java | 411 ++++++++++++++++++ .../durable/operation/InvokeOperation.java | 23 +- .../plugin/DurableExecutionPlugin.java | 7 +- .../lambda/durable/plugin/PluginRunner.java | 3 +- .../durable/plugin/PropagationMetadata.java | 4 +- .../InvokePropagationSerializationTest.java | 161 +++++++ .../operation/InvokePropagationErrorTest.java | 62 +++ 9 files changed, 718 insertions(+), 43 deletions(-) create mode 100644 otel-plugin/src/test/java/software/amazon/lambda/durable/otel/InvokePropagationIntegrationTest.java create mode 100644 sdk/src/test/java/software/amazon/lambda/durable/client/InvokePropagationSerializationTest.java create mode 100644 sdk/src/test/java/software/amazon/lambda/durable/operation/InvokePropagationErrorTest.java diff --git a/docs/advanced/propagation-metadata.md b/docs/advanced/propagation-metadata.md index 98db16349..c8ba58b94 100644 --- a/docs/advanced/propagation-metadata.md +++ b/docs/advanced/propagation-metadata.md @@ -1,39 +1,64 @@ -# Chained-invoke propagation groundwork +# Chained-invoke trace propagation -This draft prepares a plugin contract and pure OTel producers. It does **not** send propagation metadata or complete -chained-invoke propagation. Current public Lambda client models do not expose `ChainedInvokeOptions.XAmznTraceId` -or `DistributedMapOptions`. Production invoke START requests do not call this hook, and no serializer or raw HTTP -workaround is included. +For a new `CHAINED_INVOKE` operation, `InvokeOperation.startInvocation()` collects optional metadata after +`onOperationStart` has established the operation's span and before checkpoint transport serialization. The customer +payload first uses its existing serializer, so a payload failure does not trigger optional metadata collection. It writes +that header through the generated `ChainedInvokeOptions.builder().xAmznTraceId(...)` method. The wire member is the +flat optional `XAmznTraceId`; there is no nested `PropagationMetadata` structure or W3C transport field. + +Function name, tenant ID, payload, operation name, ID, parent ID and checkpoint behavior retain their existing +paths. Each operation carries its own header, including when several invokes share a checkpoint request. No +execution-global or HTTP-request trace header can substitute for those distinct calling-operation parents. + +## Plugin contract and fallback `PropagationInput` is an immutable SDK-owned snapshot with `executionArn`, `operationId`, optional `parentOperationId`, and `targetFunctionName`. `PropagationMetadata` is an immutable SDK-owned contribution with optional `xAmznTraceId`. Both types use only Java strings; generated Lambda and OTel model types do not enter the plugin interface. The typed default hook adapts a String-only producer overload, so bundled plugin layers do not reference newly added -core classes. Older cores continue to load and run those layers; new metadata collection requires a newer core, -without changing registration/provider API versions or rejecting previously working plugins. - -The optional synchronous `providePropagationMetadata` default method returns null. Existing plugins keep compiling -and linking. The dispatcher supplies the same immutable snapshot in configured order. The first non-null supported -member wins; equal values do not conflict. Different later values warn with the winning and conflicting plugin -identities and a conflict count, without logging the header. Null contributions abstain; blank headers and ordinary -hook failures are logged and skipped. Java's return type rules out asynchronous and foreign result objects. -The collector preserves the current event dispatch policy: `Exception` (including `CancellationException`) is +core classes. Older cores continue to load those layers and retain their original behavior. New collection requires +an updated core; registration/provider API versions and existing hooks are unchanged. + +The optional synchronous `providePropagationMetadata` default method returns null. The dispatcher supplies the same +immutable snapshot in configured order. The first non-null supported member wins; equal values do not conflict. +Different later values warn with the winning and conflicting plugin identities and a conflict count, without logging +the header. Null contributions abstain; blank headers and ordinary hook failures are logged and skipped. A missing +contribution omits the transport member so the backend can retain its existing inherited-header fallback. +The collector preserves Java's current event dispatch policy: `Exception` (including `CancellationException`) is contained, while `Error` propagates. A failing logging backend does not interrupt collection. -Both OTel views encode their resolved canonical trace ID, operation span ID and sampling decision as -`Root=1-<8 hex>-<24 hex>;Parent=<16 hex>;Sampled=0|1`. They require an active invocation with matching execution -ownership. An observed operation's context is authoritative, including an invocation-view continuation segment. -Before a new operation starts, the producer derives the same initial-operation ID the normal start hook will use. -It creates no spans, samples nothing again, and does not alter the existing provider, resource, upstream-parent or -fallback behavior. An explicit upstream `Sampled=0` still produces metadata with `Sampled=0`. +Both OTel views encode the resolved canonical trace ID, actual calling-operation span ID and sampling decision as +`Root=1-<8 hex>-<24 hex>;Parent=<16 hex>;Sampled=0|1`. They require active invocation ownership. The producer does not +create a span or sample again, and retains the configured provider/resource and upstream-parent behavior. Explicit +upstream `Sampled=0` is preserved. The producer can derive an initial operation ID before a start hook, but the real +invoke path collects after that hook and therefore uses the operation context it already established. + +## Replay and retry + +Replaying a stored START polls for the existing operation; replaying a terminal operation returns or raises its +stored outcome. Neither path collects metadata or sends another START. If a START checkpoint failed or was lost before +commit, a later invocation can create it again and collect again. Plugins must be deterministic and side-effect free; +this hook is not exactly-once delivery. Checkpoint batching and retries preserve each submitted operation's options. + +## Model and backend dependencies + +The SDK path and tests are implemented assuming the reviewed generated model member exists. The pinned public Lambda +model `2.55.6` currently lacks `ChainedInvokeOptions.Builder.xAmznTraceId(String)`, so compilation against that model is +expected to fail. The implementation does not hide the missing member with reflection, runtime capability checks, +serializer bypasses or raw HTTP fields. Rebase onto the published model and rerun normal tests when it is available. -Remaining work before this can implement the feature: +The design also adds `XAmznTraceId` to `DistributedMapOptions`. This Java SDK currently exposes no distributed-map +operation or START dispatch; its existing `map` and `parallel` APIs use CONTEXT operations. There is no new distributed +map/fanout/HTTP API here. The future generated model and high-level operation path must integrate their corresponding +field when introduced. -1. Supported public client models and generated serialization for `ChainedInvokeOptions.XAmznTraceId` and - `DistributedMapOptions`. -2. Backend rollout of the reviewed flat X-Ray header contract. -3. Consuming the collector on the real operation START path, including replay and failed-checkpoint integration. -4. Deployed trace topology and propagation validation. +Before release, the generated model must be published and the backend must persist/forward the per-operation header. +Then validate deployed parent/child traces for durable and ordinary Lambda targets. Keep this PR draft until those +dependencies are ready; this remains related to issue #764. Runtime inbound headers and plugin-instance lifetime are +separate changes. -Keep the implementation draft while these dependencies block completion. This is related to Java issue #764; -it does not close that issue. Runtime inbound headers and plugin-instance lifetime remain separate changes. +Tests exercise real public invoke and checkpoint paths, pending/terminal replay, failed uncommitted START recovery, +ordinary plugin fallbacks and Error propagation, sampled/unsampled OTel views, batched invokes with distinct parents, +custom payload and tenant preservation, generated-model copies, and the normal Lambda client's JSON marshaller with +only HTTP transport replaced. An isolated local model-preview fixture can test the assumed member shape, but cannot +establish that the public model or backend supports it. diff --git a/otel-plugin/README.md b/otel-plugin/README.md index aceded49f..be384df0b 100644 --- a/otel-plugin/README.md +++ b/otel-plugin/README.md @@ -40,11 +40,12 @@ If you configure your own `SdkTracerProviderBuilder`, add the OpenTelemetry SDK ``` -## Chained-invoke propagation groundwork +## Chained-invoke propagation -The SDK exposes a draft, synchronous metadata contract and both views provide pure X-Ray metadata producers. -Production invoke START integration and supported public client models are still pending; this does not enable -outbound propagation. See [the scope and remaining dependencies](../docs/advanced/propagation-metadata.md). +Both views provide the calling operation's X-Ray context to the SDK's synchronous collector. New invoke START +checkpoints carry that context in the flat optional `ChainedInvokeOptions.XAmznTraceId` field. This draft requires +the corresponding generated model and backend support; the current public model does not yet compile the new +typed setter. See [the implemented path and remaining dependencies](../docs/advanced/propagation-metadata.md). ## Quick Start using X-Ray/CloudWatch Tracing (ADOT Java Agent) diff --git a/otel-plugin/src/test/java/software/amazon/lambda/durable/otel/InvokePropagationIntegrationTest.java b/otel-plugin/src/test/java/software/amazon/lambda/durable/otel/InvokePropagationIntegrationTest.java new file mode 100644 index 000000000..5b89d7f4b --- /dev/null +++ b/otel-plugin/src/test/java/software/amazon/lambda/durable/otel/InvokePropagationIntegrationTest.java @@ -0,0 +1,411 @@ +// 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 io.opentelemetry.api.common.AttributeKey; +import io.opentelemetry.sdk.testing.exporter.InMemorySpanExporter; +import io.opentelemetry.sdk.trace.SdkTracerProvider; +import io.opentelemetry.sdk.trace.export.SimpleSpanProcessor; +import io.opentelemetry.sdk.trace.samplers.Sampler; +import java.time.Duration; +import java.time.Instant; +import java.util.ArrayList; +import java.util.LinkedHashMap; +import java.util.List; +import java.util.Map; +import java.util.concurrent.CancellationException; +import java.util.concurrent.CopyOnWriteArrayList; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicInteger; +import java.util.function.BiFunction; +import java.util.function.Consumer; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.Timeout; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.CsvSource; +import org.junit.jupiter.params.provider.EnumSource; +import org.junit.jupiter.params.provider.ValueSource; +import software.amazon.awssdk.services.lambda.model.CheckpointDurableExecutionResponse; +import software.amazon.awssdk.services.lambda.model.CheckpointUpdatedExecutionState; +import software.amazon.awssdk.services.lambda.model.ErrorObject; +import software.amazon.awssdk.services.lambda.model.ExecutionDetails; +import software.amazon.awssdk.services.lambda.model.Operation; +import software.amazon.awssdk.services.lambda.model.OperationAction; +import software.amazon.awssdk.services.lambda.model.OperationStatus; +import software.amazon.awssdk.services.lambda.model.OperationType; +import software.amazon.awssdk.services.lambda.model.OperationUpdate; +import software.amazon.lambda.durable.DurableConfig; +import software.amazon.lambda.durable.DurableContext; +import software.amazon.lambda.durable.TypeToken; +import software.amazon.lambda.durable.config.InvokeConfig; +import software.amazon.lambda.durable.config.ParallelConfig; +import software.amazon.lambda.durable.exception.UnrecoverableDurableExecutionException; +import software.amazon.lambda.durable.execution.DurableExecutor; +import software.amazon.lambda.durable.model.DurableExecutionInput; +import software.amazon.lambda.durable.model.DurableExecutionOutput; +import software.amazon.lambda.durable.model.ExecutionStatus; +import software.amazon.lambda.durable.plugin.DurableExecutionPlugin; +import software.amazon.lambda.durable.plugin.OperationInfo; +import software.amazon.lambda.durable.plugin.PropagationInput; +import software.amazon.lambda.durable.plugin.PropagationMetadata; +import software.amazon.lambda.durable.serde.JacksonSerDes; +import software.amazon.lambda.durable.serde.SerDesContext; +import software.amazon.lambda.durable.testing.local.LocalMemoryExecutionClient; +import software.amazon.lambda.durable.testing.local.OperationResult; + +/** Real public invoke/checkpoint/replay paths; requires the generated XAmznTraceId model member. */ +class InvokePropagationIntegrationTest { + private static final String TRACE = "6955b900123456789012345678901234"; + private static final String ARN = + "arn:aws:lambda:us-east-1:123456789012:function:parent/durable-execution/test/execution-id"; + private static final String HEADER = "Root=1-6955b900-123456789012345678901234;Parent=1234567890123456;Sampled=1"; + private static final String PAYLOAD = "{\"traceparent\":\"customer-owned\",\"value\":42}"; + private static final BiFunction INVOKE = + (input, context) -> context.invoke("invoke", "target:live", Map.of("value", 42), String.class); + + @Test + void committedStartAndTerminalReplayNeverRecollectOrReserialize() { + var collected = new CopyOnWriteArrayList(); + var payloadCalls = new AtomicInteger(); + var plugin = producer(collected); + var invokeConfig = InvokeConfig.builder() + .tenantId("tenant-A") + .payloadSerDes(new JacksonSerDes() { + @Override + public String serialize(Object value, SerDesContext context) { + payloadCalls.incrementAndGet(); + return PAYLOAD; + } + }) + .build(); + BiFunction handler = (input, context) -> + context.invoke("invoke", "target:live", Map.of("value", 42), String.class, invokeConfig); + try (var fixture = new Fixture(plugin)) { + fixture.client.beforeInvokeCheckpoint = update -> { + assertEquals(1, collected.size(), "Metadata must be collected before checkpoint transport"); + assertEquals(1, payloadCalls.get()); + assertEquals(HEADER, update.chainedInvokeOptions().xAmznTraceId()); + assertEquals(PAYLOAD, update.payload()); + }; + assertEquals(ExecutionStatus.PENDING, fixture.run(handler).status()); + assertEquals(ExecutionStatus.PENDING, fixture.run(handler).status()); + fixture.client.completeChainedInvoke("invoke", OperationResult.succeeded("\"child-result\"")); + assertEquals("\"child-result\"", fixture.run(handler).result()); + assertEquals("\"child-result\"", fixture.run(handler).result()); + assertEquals(1, collected.size()); + assertEquals(1, payloadCalls.get()); + var update = fixture.client.starts().get(0); + assertEquals(1, fixture.client.starts().size()); + assertEquals(HEADER, update.chainedInvokeOptions().xAmznTraceId()); + assertEquals("target:live", update.chainedInvokeOptions().functionName()); + assertEquals("tenant-A", update.chainedInvokeOptions().tenantId()); + assertEquals(PAYLOAD, update.payload()); + assertEquals("invoke", update.name()); + assertEquals(OperationType.CHAINED_INVOKE, update.type()); + assertEquals(collected.get(0).operationId(), update.id()); + assertEquals(ARN, collected.get(0).executionArn()); + assertEquals("target:live", collected.get(0).targetFunctionName()); + assertNull(collected.get(0).parentOperationId()); + } + } + + @Test + void payloadSerializationFailureDoesNotCollectOrSendStart() { + var collected = new CopyOnWriteArrayList(); + var invokeConfig = InvokeConfig.builder() + .payloadSerDes(new JacksonSerDes() { + @Override + public String serialize(Object value, SerDesContext context) { + throw new IllegalStateException("payload serialization failed"); + } + }) + .build(); + try (var fixture = new Fixture(producer(collected))) { + var result = fixture.run( + (input, context) -> context.invoke("invoke", "target:live", "payload", String.class, invokeConfig)); + assertEquals(ExecutionStatus.FAILED, result.status()); + assertEquals("payload serialization failed", result.error().errorMessage()); + assertTrue(collected.isEmpty()); + assertTrue(fixture.client.starts().isEmpty()); + } + } + + @ParameterizedTest + @EnumSource( + value = OperationStatus.class, + names = {"FAILED", "TIMED_OUT", "STOPPED"}) + void terminalFailureReplayPreservesOutcomeWithoutCollecting(OperationStatus status) { + var collected = new CopyOnWriteArrayList(); + try (var fixture = new Fixture(producer(collected))) { + assertEquals(ExecutionStatus.PENDING, fixture.run(INVOKE).status()); + fixture.client.completeChainedInvoke( + "invoke", + new OperationResult( + status, + null, + ErrorObject.builder() + .errorType("child-error") + .errorMessage("child failed") + .build())); + var first = fixture.run(INVOKE); + var replay = fixture.run(INVOKE); + assertEquals(ExecutionStatus.FAILED, first.status()); + assertEquals(first.status(), replay.status()); + assertEquals(first.error(), replay.error()); + assertEquals(1, collected.size()); + assertEquals(1, fixture.client.starts().size()); + } + } + + @ParameterizedTest + @ValueSource(strings = {"none", "legacy", "null", "empty", "blank", "exception", "cancellation"}) + void absentOrOrdinarilyFailingPluginPreservesUninstrumentedInvoke(String mode) { + DurableExecutionPlugin plugin = new DurableExecutionPlugin() { + @Override + public PropagationMetadata providePropagationMetadata(PropagationInput input) { + return switch (mode) { + case "empty" -> new PropagationMetadata(null); + case "blank" -> new PropagationMetadata(" "); + case "exception" -> throw new IllegalStateException("optional plugin failure"); + case "cancellation" -> throw new CancellationException("plugin cancellation"); + default -> null; + }; + } + }; + var plugins = mode.equals("none") + ? new DurableExecutionPlugin[0] + : new DurableExecutionPlugin[] {mode.equals("legacy") ? new DurableExecutionPlugin() {} : plugin}; + try (var fixture = new Fixture(plugins)) { + fixture.client.autoComplete = true; + assertEquals(ExecutionStatus.SUCCEEDED, fixture.run(INVOKE).status()); + var options = fixture.client.starts().get(0).chainedInvokeOptions(); + assertNull(options.xAmznTraceId()); + assertEquals("target:live", options.functionName()); + } + } + + @Test + void failedUncommittedStartCanCollectAgainWithTheSameOperationIdentity() { + var collected = new CopyOnWriteArrayList(); + try (var fixture = new Fixture(producer(collected))) { + fixture.client.failNextStart.set(true); + var error = assertThrows(UnrecoverableDurableExecutionException.class, () -> fixture.run(INVOKE)); + assertTrue(error.isRetryable()); + assertTrue(fixture.client.getAllOperations().isEmpty(), "Failed checkpoint was not committed"); + assertEquals(ExecutionStatus.PENDING, fixture.run(INVOKE).status()); + assertEquals(2, collected.size(), "Uncommitted attempts do not promise exactly-once hook calls"); + assertEquals(collected.get(0).operationId(), collected.get(1).operationId()); + assertEquals(collected.get(0).executionArn(), collected.get(1).executionArn()); + assertEquals(2, fixture.client.starts().size()); + assertEquals(fixture.client.starts().get(0), fixture.client.starts().get(1)); + } + } + + @ParameterizedTest + @CsvSource({"true,true", "true,false", "false,true", "false,false"}) + @Timeout(20) + void batchedPublicInvokesCarryTheirOwnActualOperationSpan(boolean executionView, boolean sampled) { + var exporter = InMemorySpanExporter.create(); + var builder = SdkTracerProvider.builder() + .setSampler(sampled ? Sampler.alwaysOff() : Sampler.alwaysOn()) + .addSpanProcessor(SimpleSpanProcessor.create(exporter)); + var config = OtelPluginConfig.builder() + .enableMdc(false) + .contextExtractor(() -> new ExtractedContext( + TRACE, + "1234567890123456", + sampled ? ExtractedContext.Sampling.SAMPLED : ExtractedContext.Sampling.NOT_SAMPLED)) + .build(); + DurableExecutionPlugin plugin = + executionView ? new ExecutionOtelPlugin(builder, config) : new InvocationOtelPlugin(builder, config); + var bothInvokes = new CountDownLatch(2); + var inputs = new CopyOnWriteArrayList(); + var startedOperations = new CopyOnWriteArrayList(); + DurableExecutionPlugin observer = new DurableExecutionPlugin() { + @Override + public void onOperationStart(OperationInfo info) { + startedOperations.add(info); + } + + @Override + public PropagationMetadata providePropagationMetadata(PropagationInput input) { + inputs.add(input); + bothInvokes.countDown(); + try { + assertTrue(bothInvokes.await(5, TimeUnit.SECONDS)); + } catch (InterruptedException error) { + Thread.currentThread().interrupt(); + throw new AssertionError(error); + } + return null; + } + }; + try (var fixture = new Fixture(plugin, observer)) { + fixture.client.autoComplete = true; + var result = fixture.run((input, context) -> { + try (var parallel = context.parallel( + "fanout", ParallelConfig.builder().maxConcurrency(2).build())) { + parallel.branch( + "a", + String.class, + child -> child.invoke("invoke-a", "target-a:live", "payload-a", String.class)); + parallel.branch( + "b", + String.class, + child -> child.invoke("invoke-b", "target-b:live", "payload-b", String.class)); + } + return "done"; + }); + assertEquals(ExecutionStatus.SUCCEEDED, result.status()); + assertEquals(2, inputs.size()); + var starts = fixture.client.starts(); + assertEquals(2, starts.size()); + assertTrue(fixture.client.batches.stream() + .anyMatch(batch -> batch.stream() + .filter(InvokePropagationIntegrationTest::isInvokeStart) + .count() + == 2)); + assertNotEquals(starts.get(0).id(), starts.get(1).id()); + assertNotEquals( + starts.get(0).chainedInvokeOptions().xAmznTraceId(), + starts.get(1).chainedInvokeOptions().xAmznTraceId()); + for (var update : starts) { + var metadata = inputs.stream() + .filter(item -> item.operationId().equals(update.id())) + .findFirst() + .orElseThrow(); + assertEquals(update.parentId(), metadata.parentOperationId()); + assertNotNull(metadata.parentOperationId()); + assertEquals(update.chainedInvokeOptions().functionName(), metadata.targetFunctionName()); + var spanId = new DeterministicIdGenerator().generateSpanIdForOperation(ARN, update.id()); + assertEquals( + "Root=1-6955b900-123456789012345678901234;Parent=" + spanId + ";Sampled=" + + (sampled ? "1" : "0"), + update.chainedInvokeOptions().xAmznTraceId()); + if (sampled) { + var span = exporter.getFinishedSpanItems().stream() + .filter(item -> update.id() + .equals(item.getAttributes() + .get(AttributeKey.stringKey("durable.operation.id"))) + && update.name().equals(item.getName())) + .findFirst() + .orElseThrow(); + assertEquals(span.getSpanId(), spanId); + assertEquals(TRACE, span.getTraceId()); + } + } + assertEquals( + sampled ? startedOperations.size() + 2 : 0, + exporter.getFinishedSpanItems().size(), + "No extra span is created to obtain propagation metadata"); + } + } + + private static DurableExecutionPlugin producer(List inputs) { + return new DurableExecutionPlugin() { + @Override + public PropagationMetadata providePropagationMetadata(PropagationInput input) { + inputs.add(input); + return new PropagationMetadata(HEADER); + } + }; + } + + private static boolean isInvokeStart(OperationUpdate update) { + return update.type() == OperationType.CHAINED_INVOKE && update.action() == OperationAction.START; + } + + private static final class RecordingClient extends LocalMemoryExecutionClient { + final List> batches = new CopyOnWriteArrayList<>(); + final AtomicBoolean failNextStart = new AtomicBoolean(); + boolean autoComplete; + Consumer beforeInvokeCheckpoint = update -> {}; + + @Override + public CheckpointDurableExecutionResponse checkpoint(String arn, String token, List updates) { + batches.add(List.copyOf(updates)); + var starts = updates.stream() + .filter(InvokePropagationIntegrationTest::isInvokeStart) + .toList(); + starts.forEach(beforeInvokeCheckpoint); + if (!starts.isEmpty() && failNextStart.compareAndSet(true, false)) { + throw new UnrecoverableDurableExecutionException( + ErrorObject.builder() + .errorType("CheckpointUnavailable") + .errorMessage("simulated uncommitted START") + .build(), + true); + } + var response = super.checkpoint(arn, token, updates); + if (autoComplete && !starts.isEmpty()) { + starts.forEach( + update -> completeChainedInvoke(update.name(), OperationResult.succeeded("\"completed\""))); + // Consume the simulator's newly completed results once. Returning all stored operations leaves + // these updates pending and spuriously delivers the same terminal event in the next checkpoint. + var completed = super.checkpoint(arn, response.checkpointToken(), List.of()); + var changed = new LinkedHashMap(); + response.newExecutionState().operations().forEach(operation -> changed.put(operation.id(), operation)); + completed.newExecutionState().operations().forEach(operation -> changed.put(operation.id(), operation)); + return completed.toBuilder() + .newExecutionState(CheckpointUpdatedExecutionState.builder() + .operations(changed.values()) + .build()) + .build(); + } + return response; + } + + List starts() { + return batches.stream() + .flatMap(List::stream) + .filter(InvokePropagationIntegrationTest::isInvokeStart) + .toList(); + } + } + + private static final class Fixture implements AutoCloseable { + final RecordingClient client = new RecordingClient(); + final ExecutorService executor = Executors.newCachedThreadPool(); + final DurableConfig config; + + Fixture(DurableExecutionPlugin... plugins) { + config = DurableConfig.builder() + .withDurableExecutionClient(client) + .withExecutorService(executor) + .withCheckpointDelay(Duration.ofMillis(250)) + .withPlugins(plugins) + .build(); + } + + DurableExecutionOutput run(BiFunction handler) { + var operations = new ArrayList<>(List.of(Operation.builder() + .id("execution-id") + .type(OperationType.EXECUTION) + .status(OperationStatus.STARTED) + .startTimestamp(Instant.parse("2026-10-05T00:00:00Z")) + .executionDetails( + ExecutionDetails.builder().inputPayload("\"input\"").build()) + .build())); + operations.addAll(client.getAllOperations()); + var input = new DurableExecutionInput( + ARN, + "token", + CheckpointUpdatedExecutionState.builder() + .operations(operations) + .build(), + client.getUpdatedOperationIdsSinceLastInvocation()); + return DurableExecutor.execute(input, null, TypeToken.get(String.class), handler, config); + } + + @Override + public void close() { + executor.shutdownNow(); + } + } +} diff --git a/sdk/src/main/java/software/amazon/lambda/durable/operation/InvokeOperation.java b/sdk/src/main/java/software/amazon/lambda/durable/operation/InvokeOperation.java index 52d365c30..bf4edb74a 100644 --- a/sdk/src/main/java/software/amazon/lambda/durable/operation/InvokeOperation.java +++ b/sdk/src/main/java/software/amazon/lambda/durable/operation/InvokeOperation.java @@ -14,6 +14,7 @@ import software.amazon.lambda.durable.exception.InvokeStoppedException; import software.amazon.lambda.durable.exception.InvokeTimedOutException; import software.amazon.lambda.durable.model.OperationIdentifier; +import software.amazon.lambda.durable.plugin.PropagationInput; import software.amazon.lambda.durable.serde.SerDes; /** @@ -64,13 +65,25 @@ protected void replay(Operation existing) { } private void startInvocation() { + var serializedPayload = payloadSerDes.serialize(this.payload, serDesContext("invoke-payload")); + var options = ChainedInvokeOptions.builder().functionName(functionName).tenantId(invokeConfig.tenantId()); + var plugins = getContext().getDurableConfig().getPluginRunner(); + if (!plugins.isEmpty()) { + // onOperationStart has already established the calling operation's span. Only a new START + // collects metadata; replay consumes the persisted outcome without recomputing it. + var metadata = plugins.providePropagationMetadata(new PropagationInput( + executionManager.getDurableExecutionArn(), + getOperationId(), + getContext().getParentId(), + functionName)); + if (metadata != null && metadata.xAmznTraceId() != null) { + options.xAmznTraceId(metadata.xAmznTraceId()); + } + } var update = OperationUpdate.builder() .action(OperationAction.START) - .chainedInvokeOptions(ChainedInvokeOptions.builder() - .functionName(functionName) - .tenantId(invokeConfig.tenantId()) - .build()) - .payload(payloadSerDes.serialize(this.payload, serDesContext("invoke-payload"))); + .chainedInvokeOptions(options.build()) + .payload(serializedPayload); sendOperationUpdate(update); } diff --git a/sdk/src/main/java/software/amazon/lambda/durable/plugin/DurableExecutionPlugin.java b/sdk/src/main/java/software/amazon/lambda/durable/plugin/DurableExecutionPlugin.java index fa5603c5c..0c28e6a7d 100644 --- a/sdk/src/main/java/software/amazon/lambda/durable/plugin/DurableExecutionPlugin.java +++ b/sdk/src/main/java/software/amazon/lambda/durable/plugin/DurableExecutionPlugin.java @@ -16,9 +16,10 @@ public interface DurableExecutionPlugin { /** - * Optionally supplies operation-level trace propagation metadata synchronously. Return null to abstain. This - * groundwork hook is not yet called by production invoke START requests; client-model support and backend rollout - * are pending. Implementations must not mutate execution state or create side effects to produce metadata. + * Optionally supplies operation-level trace propagation metadata synchronously. Return null to abstain. This hook + * runs when creating a new invoke START checkpoint after the operation-start hook. Replaying a checkpointed START + * or terminal result does not call it; an uncommitted START may call it again on a later invocation. Produce + * metadata deterministically from the input without mutating execution state or creating side effects. */ default PropagationMetadata providePropagationMetadata(PropagationInput input) { var header = providePropagationMetadata( diff --git a/sdk/src/main/java/software/amazon/lambda/durable/plugin/PluginRunner.java b/sdk/src/main/java/software/amazon/lambda/durable/plugin/PluginRunner.java index d5f18d8b2..d39ab02b0 100644 --- a/sdk/src/main/java/software/amazon/lambda/durable/plugin/PluginRunner.java +++ b/sdk/src/main/java/software/amazon/lambda/durable/plugin/PluginRunner.java @@ -47,7 +47,8 @@ public List getPlugins() { /** * Collects supported metadata in configured order on the caller's thread. First non-null member wins; matching * later values are harmless. Ordinary plugin/invalid-result failures are logged and skipped, matching the event - * hooks' Exception containment policy; Errors continue to propagate. No production START path calls this yet. + * hooks' Exception containment policy; Errors continue to propagate. The core calls this while creating a new + * invoke START checkpoint, not while replaying an existing operation. */ public PropagationMetadata providePropagationMetadata(PropagationInput input) { requireNonNull(input, "input"); diff --git a/sdk/src/main/java/software/amazon/lambda/durable/plugin/PropagationMetadata.java b/sdk/src/main/java/software/amazon/lambda/durable/plugin/PropagationMetadata.java index 1b5fededd..8eeeb3a2d 100644 --- a/sdk/src/main/java/software/amazon/lambda/durable/plugin/PropagationMetadata.java +++ b/sdk/src/main/java/software/amazon/lambda/durable/plugin/PropagationMetadata.java @@ -3,8 +3,8 @@ package software.amazon.lambda.durable.plugin; /** - * Immutable SDK-owned propagation contribution. This is not a generated Lambda request model and is not currently - * serialized on the production invoke START path. A null header means the plugin has no contribution. + * Immutable SDK-owned propagation contribution, separate from generated Lambda request models. The core maps its + * optional header to the invoke START checkpoint. A null header means the plugin has no contribution. */ public final class PropagationMetadata { private final String xAmznTraceId; diff --git a/sdk/src/test/java/software/amazon/lambda/durable/client/InvokePropagationSerializationTest.java b/sdk/src/test/java/software/amazon/lambda/durable/client/InvokePropagationSerializationTest.java new file mode 100644 index 000000000..16e9d081a --- /dev/null +++ b/sdk/src/test/java/software/amazon/lambda/durable/client/InvokePropagationSerializationTest.java @@ -0,0 +1,161 @@ +// Copyright Amazon.com, Inc. or its affiliates. All Rights Reserved. +// SPDX-License-Identifier: Apache-2.0 +package software.amazon.lambda.durable.client; + +import static org.junit.jupiter.api.Assertions.*; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.*; + +import com.fasterxml.jackson.databind.JsonNode; +import com.fasterxml.jackson.databind.ObjectMapper; +import java.io.ByteArrayInputStream; +import java.io.IOException; +import java.nio.charset.StandardCharsets; +import java.util.ArrayList; +import java.util.List; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.ValueSource; +import software.amazon.awssdk.auth.credentials.AwsBasicCredentials; +import software.amazon.awssdk.auth.credentials.StaticCredentialsProvider; +import software.amazon.awssdk.core.SdkRequest; +import software.amazon.awssdk.core.interceptor.Context; +import software.amazon.awssdk.core.interceptor.ExecutionAttributes; +import software.amazon.awssdk.core.interceptor.ExecutionInterceptor; +import software.amazon.awssdk.http.AbortableInputStream; +import software.amazon.awssdk.http.ExecutableHttpRequest; +import software.amazon.awssdk.http.HttpExecuteRequest; +import software.amazon.awssdk.http.HttpExecuteResponse; +import software.amazon.awssdk.http.SdkHttpClient; +import software.amazon.awssdk.http.SdkHttpResponse; +import software.amazon.awssdk.regions.Region; +import software.amazon.awssdk.services.lambda.LambdaClient; +import software.amazon.awssdk.services.lambda.model.ChainedInvokeOptions; +import software.amazon.awssdk.services.lambda.model.CheckpointDurableExecutionRequest; +import software.amazon.awssdk.services.lambda.model.OperationAction; +import software.amazon.awssdk.services.lambda.model.OperationType; +import software.amazon.awssdk.services.lambda.model.OperationUpdate; + +/** Uses the real generated Lambda protocol marshaller; only the HTTP exchange is replaced. */ +class InvokePropagationSerializationTest { + private static final String ARN = "arn:aws:lambda:us-east-1:123456789012:function:parent/durable-execution/test/id"; + private static final String PAYLOAD = "{\"traceparent\":\"customer-owned\",\"nested\":{\"value\":42}}"; + + @ParameterizedTest + @ValueSource(strings = {"1", "0", "absent"}) + void checkpointSerializationPreservesFlatPerOperationHeadersThroughModelCopy(String sampled) throws Exception { + var updates = List.of(update("a", sampled), update("b", sampled)); + // Exercise generated immutable-model copy paths before the real client serializer. + var copied = updates.stream() + .map(item -> item.toBuilder() + .chainedInvokeOptions( + item.chainedInvokeOptions().toBuilder().build()) + .build()) + .toList(); + var requests = new ArrayList(); + var mapper = new ObjectMapper(); + var http = mock(SdkHttpClient.class); + when(http.prepareRequest(any())).thenAnswer(call -> { + HttpExecuteRequest request = call.getArgument(0); + try (var body = request.contentStreamProvider().orElseThrow().newStream()) { + requests.add(mapper.readTree(body)); + } + return new ExecutableHttpRequest() { + @Override + public HttpExecuteResponse call() throws IOException { + var body = new ByteArrayInputStream( + "{\"CheckpointToken\":\"next-token\"}".getBytes(StandardCharsets.UTF_8)); + return HttpExecuteResponse.builder() + .response(SdkHttpResponse.builder() + .statusCode(200) + .putHeader("Content-Type", "application/json") + .build()) + .responseBody(AbortableInputStream.create(body)) + .build(); + } + + @Override + public void abort() {} + }; + }); + try (var lambda = LambdaClient.builder() + .region(Region.US_EAST_1) + .httpClient(http) + .overrideConfiguration(configuration -> + configuration.addExecutionInterceptor(new ExecutionInterceptor() { + @Override + public SdkRequest modifyRequest( + Context.ModifyRequest context, ExecutionAttributes attributes) { + // Stabilize only this fixture's automatic idempotency field, preserving the full wire + // comparison. + return ((CheckpointDurableExecutionRequest) context.request()) + .toBuilder() + .clientToken("propagation-serialization-fixture") + .build(); + } + })) + .credentialsProvider( + StaticCredentialsProvider.create(AwsBasicCredentials.create("test-key", "test-secret"))) + .build()) { + var client = new LambdaDurableFunctionsClient(lambda); + assertEquals("next-token", client.checkpoint(ARN, "token", updates).checkpointToken()); + assertEquals("next-token", client.checkpoint(ARN, "token", copied).checkpointToken()); + } + assertEquals(2, requests.size()); + assertEquals(requests.get(0), requests.get(1), "Model copies must preserve the complete checkpoint payload"); + assertEquals( + "propagation-serialization-fixture", + requests.get(0).get("ClientToken").asText()); + var serialized = requests.get(0).get("Updates"); + assertEquals(2, serialized.size()); + for (var index = 0; index < updates.size(); index++) { + var original = updates.get(index); + var wire = serialized.get(index); + assertEquals(original.id(), wire.get("Id").asText()); + assertEquals(original.parentId(), wire.get("ParentId").asText()); + assertEquals(original.name(), wire.get("Name").asText()); + assertEquals("CHAINED_INVOKE", wire.get("Type").asText()); + assertEquals("START", wire.get("Action").asText()); + assertEquals(PAYLOAD, wire.get("Payload").asText()); + var options = wire.get("ChainedInvokeOptions"); + assertEquals( + original.chainedInvokeOptions().functionName(), + options.get("FunctionName").asText()); + assertEquals( + original.chainedInvokeOptions().tenantId(), + options.get("TenantId").asText()); + assertFalse(options.has("PropagationMetadata"), "The model contract is flat, not the older nested sketch"); + if (sampled.equals("absent")) { + assertFalse(options.has("XAmznTraceId"), "Absent contribution must omit the field, not send null"); + } else { + assertEquals( + original.chainedInvokeOptions().xAmznTraceId(), + options.get("XAmznTraceId").asText()); + } + } + if (!sampled.equals("absent")) { + assertNotEquals( + serialized.get(0).get("ChainedInvokeOptions").get("XAmznTraceId"), + serialized.get(1).get("ChainedInvokeOptions").get("XAmznTraceId")); + } + } + + private static OperationUpdate update(String suffix, String sampled) { + var options = ChainedInvokeOptions.builder() + .functionName("target-" + suffix + ":live") + .tenantId("tenant-" + suffix); + if (!sampled.equals("absent")) { + options.xAmznTraceId( + "Root=1-6955b900-123456789012345678901234;Parent=" + suffix.repeat(16) + ";Sampled=" + sampled); + } + return OperationUpdate.builder() + .id("invoke-" + suffix) + .parentId("context-" + suffix) + .name("call-" + suffix) + .type(OperationType.CHAINED_INVOKE) + .subType("ChainedInvoke") + .action(OperationAction.START) + .payload(PAYLOAD) + .chainedInvokeOptions(options.build()) + .build(); + } +} diff --git a/sdk/src/test/java/software/amazon/lambda/durable/operation/InvokePropagationErrorTest.java b/sdk/src/test/java/software/amazon/lambda/durable/operation/InvokePropagationErrorTest.java new file mode 100644 index 000000000..96e945eb7 --- /dev/null +++ b/sdk/src/test/java/software/amazon/lambda/durable/operation/InvokePropagationErrorTest.java @@ -0,0 +1,62 @@ +// Copyright Amazon.com, Inc. or its affiliates. All Rights Reserved. +// SPDX-License-Identifier: Apache-2.0 +package software.amazon.lambda.durable.operation; + +import static org.junit.jupiter.api.Assertions.*; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.*; + +import java.util.List; +import java.util.concurrent.atomic.AtomicBoolean; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.ValueSource; +import software.amazon.lambda.durable.DurableConfig; +import software.amazon.lambda.durable.TypeToken; +import software.amazon.lambda.durable.config.InvokeConfig; +import software.amazon.lambda.durable.context.DurableContextImpl; +import software.amazon.lambda.durable.execution.ExecutionManager; +import software.amazon.lambda.durable.model.OperationIdentifier; +import software.amazon.lambda.durable.model.OperationSubType; +import software.amazon.lambda.durable.plugin.DurableExecutionPlugin; +import software.amazon.lambda.durable.plugin.OperationInfo; +import software.amazon.lambda.durable.plugin.PluginRunner; +import software.amazon.lambda.durable.plugin.PropagationInput; +import software.amazon.lambda.durable.plugin.PropagationMetadata; +import software.amazon.lambda.durable.serde.JacksonSerDes; + +class InvokePropagationErrorTest { + @ParameterizedTest + @ValueSource(booleans = {false, true}) + void collectorErrorsRetainIdentityAndPreventCheckpoint(boolean linkageError) { + Error failure = linkageError ? new NoSuchMethodError("plugin linkage") : new AssertionError("plugin error"); + var started = new AtomicBoolean(); + DurableExecutionPlugin plugin = new DurableExecutionPlugin() { + @Override + public void onOperationStart(OperationInfo info) { + started.set(true); + } + + @Override + public PropagationMetadata providePropagationMetadata(PropagationInput input) { + assertTrue(started.get(), "Operation context is established before collection"); + throw failure; + } + }; + var manager = mock(ExecutionManager.class); + when(manager.getDurableExecutionArn()).thenReturn("execution-arn"); + var config = mock(DurableConfig.class); + when(config.getPluginRunner()).thenReturn(new PluginRunner(List.of(plugin))); + var context = mock(DurableContextImpl.class); + when(context.getExecutionManager()).thenReturn(manager); + when(context.getDurableConfig()).thenReturn(config); + var operation = new InvokeOperation<>( + OperationIdentifier.of("invoke-id", "invoke", OperationSubType.CHAINED_INVOKE), + "target", + "payload", + TypeToken.get(String.class), + InvokeConfig.builder().serDes(new JacksonSerDes()).build(), + context); + assertSame(failure, assertThrows(Error.class, operation::execute)); + verify(manager, never()).sendOperationUpdate(any()); + } +} From b744341fa63d0d5445d62ac1a497ffb21aca4d4f Mon Sep 17 00:00:00 2001 From: Frank Chen <65260095+zhongkechen@users.noreply.github.com> Date: Mon, 5 Oct 2026 22:57:23 -0700 Subject: [PATCH 3/7] fix: account for chained invoke options in checkpoint batch sizes --- docs/advanced/propagation-metadata.md | 12 +- .../durable/execution/CheckpointManager.java | 29 +++- .../execution/InvokeCheckpointSizeTest.java | 156 ++++++++++++++++++ 3 files changed, 191 insertions(+), 6 deletions(-) create mode 100644 sdk/src/test/java/software/amazon/lambda/durable/execution/InvokeCheckpointSizeTest.java diff --git a/docs/advanced/propagation-metadata.md b/docs/advanced/propagation-metadata.md index c8ba58b94..30b000144 100644 --- a/docs/advanced/propagation-metadata.md +++ b/docs/advanced/propagation-metadata.md @@ -43,9 +43,10 @@ this hook is not exactly-once delivery. Checkpoint batching and retries preserve ## Model and backend dependencies The SDK path and tests are implemented assuming the reviewed generated model member exists. The pinned public Lambda -model `2.55.6` currently lacks `ChainedInvokeOptions.Builder.xAmznTraceId(String)`, so compilation against that model is -expected to fail. The implementation does not hide the missing member with reflection, runtime capability checks, -serializer bypasses or raw HTTP fields. Rebase onto the published model and rerun normal tests when it is available. +model `2.55.11` currently lacks `ChainedInvokeOptions.XAmznTraceId` (both its builder setter and getter), so compilation +against that model is expected to fail. The implementation does not hide the missing member with reflection, runtime +capability checks, serializer bypasses or raw HTTP fields. Rebase onto the published model and rerun normal tests when +it is available. The design also adds `XAmznTraceId` to `DistributedMapOptions`. This Java SDK currently exposes no distributed-map operation or START dispatch; its existing `map` and `parallel` APIs use CONTEXT operations. There is no new distributed @@ -60,5 +61,6 @@ separate changes. Tests exercise real public invoke and checkpoint paths, pending/terminal replay, failed uncommitted START recovery, ordinary plugin fallbacks and Error propagation, sampled/unsampled OTel views, batched invokes with distinct parents, custom payload and tenant preservation, generated-model copies, and the normal Lambda client's JSON marshaller with -only HTTP transport replaced. An isolated local model-preview fixture can test the assumed member shape, but cannot -establish that the public model or backend supports it. +only HTTP transport replaced. Batch-boundary tests include the encoded options and trace header in the 750 KiB size +budget, including JSON escaping and UTF-8. An isolated local model-preview fixture can test the assumed member shape, +but cannot establish that the public model or backend supports it. diff --git a/sdk/src/main/java/software/amazon/lambda/durable/execution/CheckpointManager.java b/sdk/src/main/java/software/amazon/lambda/durable/execution/CheckpointManager.java index 34f44134c..1b833b560 100644 --- a/sdk/src/main/java/software/amazon/lambda/durable/execution/CheckpointManager.java +++ b/sdk/src/main/java/software/amazon/lambda/durable/execution/CheckpointManager.java @@ -266,10 +266,37 @@ private static int estimateSize(OperationUpdate update) { if (update == null) { return 0; } - return update.id().length() + long size = (long) update.id().length() + update.type().toString().length() + update.action().toString().length() + (update.payload() != null ? update.payload().length() : 0) + 100; + var options = update.chainedInvokeOptions(); + if (options != null) { + size += ",\"ChainedInvokeOptions\":{}".length(); + size += jsonStringFieldSize("FunctionName", options.functionName()); + size += jsonStringFieldSize("TenantId", options.tenantId()); + size += jsonStringFieldSize("XAmznTraceId", options.xAmznTraceId()); + } + return (int) Math.min(size, Integer.MAX_VALUE); + } + + /** Include UTF-8 and JSON escaping; conservatively allow a comma before every optional field. */ + private static long jsonStringFieldSize(String name, String value) { + if (value == null) return 0; + long size = name.length() + 6L; // comma, quoted key, colon, and quoted value + for (int index = 0; index < value.length(); index++) { + char character = value.charAt(index); + if (character < 0x20 || character >= 0x80) { + // Allow the longest JSON escape per UTF-16 unit. The Lambda marshaller escapes surrogate + // pairs as two six-byte escapes, rather than emitting a four-byte UTF-8 code point. + size += 6; + } else if (character == '"' || character == '\\') { + size += 2; + } else { + size++; + } + } + return size; } } diff --git a/sdk/src/test/java/software/amazon/lambda/durable/execution/InvokeCheckpointSizeTest.java b/sdk/src/test/java/software/amazon/lambda/durable/execution/InvokeCheckpointSizeTest.java new file mode 100644 index 000000000..0ba9484f6 --- /dev/null +++ b/sdk/src/test/java/software/amazon/lambda/durable/execution/InvokeCheckpointSizeTest.java @@ -0,0 +1,156 @@ +// Copyright Amazon.com, Inc. or its affiliates. All Rights Reserved. +// SPDX-License-Identifier: Apache-2.0 +package software.amazon.lambda.durable.execution; + +import static org.junit.jupiter.api.Assertions.*; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.*; + +import com.fasterxml.jackson.databind.ObjectMapper; +import java.io.ByteArrayInputStream; +import java.io.IOException; +import java.nio.charset.StandardCharsets; +import java.time.Duration; +import java.util.List; +import java.util.concurrent.CopyOnWriteArrayList; +import java.util.concurrent.TimeUnit; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.ValueSource; +import software.amazon.awssdk.auth.credentials.AwsBasicCredentials; +import software.amazon.awssdk.auth.credentials.StaticCredentialsProvider; +import software.amazon.awssdk.core.SdkRequest; +import software.amazon.awssdk.core.interceptor.Context; +import software.amazon.awssdk.core.interceptor.ExecutionAttributes; +import software.amazon.awssdk.core.interceptor.ExecutionInterceptor; +import software.amazon.awssdk.http.AbortableInputStream; +import software.amazon.awssdk.http.ExecutableHttpRequest; +import software.amazon.awssdk.http.HttpExecuteRequest; +import software.amazon.awssdk.http.HttpExecuteResponse; +import software.amazon.awssdk.http.SdkHttpClient; +import software.amazon.awssdk.http.SdkHttpResponse; +import software.amazon.awssdk.regions.Region; +import software.amazon.awssdk.services.lambda.LambdaClient; +import software.amazon.awssdk.services.lambda.model.ChainedInvokeOptions; +import software.amazon.awssdk.services.lambda.model.CheckpointDurableExecutionRequest; +import software.amazon.awssdk.services.lambda.model.OperationAction; +import software.amazon.awssdk.services.lambda.model.OperationType; +import software.amazon.awssdk.services.lambda.model.OperationUpdate; +import software.amazon.lambda.durable.DurableConfig; +import software.amazon.lambda.durable.client.LambdaDurableFunctionsClient; + +/** Measures the real generated protocol body; only HTTP exchange is replaced. */ +class InvokeCheckpointSizeTest { + private static final int BUDGET = 750 * 1024; + private static final String ARN = "arn:aws:lambda:us-east-1:123456789012:function:parent/durable-execution/test/id"; + private static final String HEADER = "Root=1-6955b900-123456789012345678901234;Parent=1234567890123456;Sampled=1"; + + @ParameterizedTest + @ValueSource(booleans = {false, true}) + void traceOptionsSplitAnOtherwiseFittingBatchWithoutLosingUpdates(boolean escaped) throws Exception { + var header = escaped ? HEADER + ";extension=" + "\\\"\n\u0001\u00e9\ud83d\ude80".repeat(128) : HEADER; + try (var wire = new Wire()) { + var empty = List.of(update("a", "", null), update("b", "", null)); + var overhead = wire.body(empty).length; + var headerBytes = wire.body(List.of(update("a", "", header), update("b", "", header))).length - overhead; + var payload = "x".repeat((BUDGET - overhead - headerBytes + 40) / 2); + var plain = List.of(update("a", payload, null), update("b", payload, null)); + assertTrue(wire.body(plain).length <= BUDGET, "control batch fits before adding trace headers"); + var updates = List.of(update("a", payload, header), update("b", payload, header)); + assertTrue(wire.body(updates).length > BUDGET, "headers cross the real JSON byte boundary"); + wire.requests.clear(); + var config = DurableConfig.builder() + .withDurableExecutionClient(wire.client) + .withCheckpointDelay(Duration.ofHours(1)) + .build(); + var manager = new CheckpointManager(config, ARN, "token", operations -> {}); + var first = manager.checkpoint(updates.get(0)); + var second = manager.checkpoint(updates.get(1)); + manager.shutdown(); + first.get(5, TimeUnit.SECONDS); + second.get(5, TimeUnit.SECONDS); + assertEquals(2, wire.requests.size(), "the added encoded options must force two batches"); + var mapper = new ObjectMapper(); + for (var index = 0; index < wire.requests.size(); index++) { + var bytes = wire.requests.get(index); + assertTrue(bytes.length <= BUDGET, "each actual wire request stays within the batching budget"); + var items = mapper.readTree(bytes).get("Updates"); + assertEquals(1, items.size()); + var item = items.get(0); + assertEquals(updates.get(index).id(), item.get("Id").asText()); + assertEquals(updates.get(index).payload(), item.get("Payload").asText()); + assertEquals( + header, + item.get("ChainedInvokeOptions").get("XAmznTraceId").asText()); + } + } + } + + private static OperationUpdate update(String id, String payload, String header) { + var options = ChainedInvokeOptions.builder().functionName("target:live").tenantId("tenant"); + if (header != null) options.xAmznTraceId(header); + return OperationUpdate.builder() + .id(id) + .type(OperationType.CHAINED_INVOKE) + .action(OperationAction.START) + .payload("\"" + payload + "\"") + .chainedInvokeOptions(options.build()) + .build(); + } + + private static final class Wire implements AutoCloseable { + final List requests = new CopyOnWriteArrayList<>(); + final LambdaClient lambda; + final LambdaDurableFunctionsClient client; + + Wire() throws IOException { + var http = mock(SdkHttpClient.class); + when(http.prepareRequest(any())).thenAnswer(invocation -> { + HttpExecuteRequest request = invocation.getArgument(0); + try (var body = request.contentStreamProvider().orElseThrow().newStream()) { + requests.add(body.readAllBytes()); + } + return new ExecutableHttpRequest() { + @Override + public HttpExecuteResponse call() { + return HttpExecuteResponse.builder() + .response(SdkHttpResponse.builder() + .statusCode(200) + .putHeader("Content-Type", "application/json") + .build()) + .responseBody(AbortableInputStream.create(new ByteArrayInputStream( + "{\"CheckpointToken\":\"token\"}".getBytes(StandardCharsets.UTF_8)))) + .build(); + } + + @Override + public void abort() {} + }; + }); + lambda = LambdaClient.builder() + .region(Region.US_EAST_1) + .httpClient(http) + .credentialsProvider(StaticCredentialsProvider.create(AwsBasicCredentials.create("test", "test"))) + .overrideConfiguration(config -> config.addExecutionInterceptor(new ExecutionInterceptor() { + @Override + public SdkRequest modifyRequest(Context.ModifyRequest context, ExecutionAttributes attributes) { + return ((CheckpointDurableExecutionRequest) context.request()) + .toBuilder() + .clientToken("batch-size-fixture") + .build(); + } + })) + .build(); + client = new LambdaDurableFunctionsClient(lambda); + } + + byte[] body(List updates) { + client.checkpoint(ARN, "token", updates); + return requests.get(requests.size() - 1); + } + + @Override + public void close() { + lambda.close(); + } + } +} From 8196cd1a1e40da11a011a13c67237b8974a01bcd Mon Sep 17 00:00:00 2001 From: Frank Chen Date: Tue, 6 Oct 2026 22:42:23 +0000 Subject: [PATCH 4/7] Preserve pending conformance and E2E CI runs --- .github/workflows/conformance-tests.yml | 1 + .github/workflows/e2e-tests.yml | 1 + 2 files changed, 2 insertions(+) diff --git a/.github/workflows/conformance-tests.yml b/.github/workflows/conformance-tests.yml index 06baa03f5..c066bf277 100644 --- a/.github/workflows/conformance-tests.yml +++ b/.github/workflows/conformance-tests.yml @@ -25,6 +25,7 @@ concurrency: # same stack -- mirrors e2e-tests.yml. group: conformance-tests cancel-in-progress: false + queue: max permissions: contents: read diff --git a/.github/workflows/e2e-tests.yml b/.github/workflows/e2e-tests.yml index e57070144..5dd246d57 100644 --- a/.github/workflows/e2e-tests.yml +++ b/.github/workflows/e2e-tests.yml @@ -26,6 +26,7 @@ on: concurrency: group: e2e-tests cancel-in-progress: false + queue: max # permission can be added at job level or workflow level permissions: From 5f83d4777b60ff8916f928cd47228cc6b7f402ca Mon Sep 17 00:00:00 2001 From: Frank Chen <65260095+zhongkechen@users.noreply.github.com> Date: Tue, 6 Oct 2026 17:39:15 -0700 Subject: [PATCH 5/7] ci: serialize revised OTel callers without replacing pending work --- .github/workflows/otel-conformance-tests.yml | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/.github/workflows/otel-conformance-tests.yml b/.github/workflows/otel-conformance-tests.yml index 68b3bc3f2..43a3e693f 100644 --- a/.github/workflows/otel-conformance-tests.yml +++ b/.github/workflows/otel-conformance-tests.yml @@ -53,6 +53,11 @@ on: permissions: {} +concurrency: + group: otel-conformance-tests + cancel-in-progress: false + queue: max + jobs: opentelemetry: # 1. Run for non-PR events, such as scheduled runs and manual invocations From 7b83565adffdf0045d309d414f21d3965e25fda8 Mon Sep 17 00:00:00 2001 From: Frank Chen Date: Wed, 7 Oct 2026 01:34:41 +0000 Subject: [PATCH 6/7] ci: run OTel conformance on validated CodeBuild Java 21 --- .github/workflows/otel-conformance-tests.yml | 33 +++++++++++++++++--- 1 file changed, 29 insertions(+), 4 deletions(-) diff --git a/.github/workflows/otel-conformance-tests.yml b/.github/workflows/otel-conformance-tests.yml index 43a3e693f..4f948b298 100644 --- a/.github/workflows/otel-conformance-tests.yml +++ b/.github/workflows/otel-conformance-tests.yml @@ -67,8 +67,9 @@ jobs: actions: write contents: read id-token: write - uses: aws/aws-durable-execution-conformance-tests/.github/workflows/opentelemetry-orchestrator.yml@43872f8d5dd917fe3ee6d01a6c953c1ff67c5fe4 + uses: aws/aws-durable-execution-conformance-tests/.github/workflows/opentelemetry-orchestrator.yml@a66037abbbfa55fde97f714e30f0bc262edefd63 with: + runs_on: ${{ github.event_name == 'pull_request' && github.event.pull_request.head.repo.full_name != github.repository && 'ubuntu-latest' || format('codebuild-github-actions-runner-{0}-{1}', github.run_id, github.run_attempt) }} language: java resource_prefix: j sdk_repository: aws/aws-durable-execution-sdk-java @@ -80,13 +81,37 @@ jobs: # checked out (.build/durable-sdk). examples_dir: .build/durable-sdk/conformance-tests-otel setup_command: | - if [ -z "${JAVA_HOME_21_X64:-}" ]; then - echo "The runner does not provide Java 21" + set -euo pipefail + if [ -n "${JAVA_HOME_21_X64:-}" ]; then + export JAVA_HOME="$JAVA_HOME_21_X64" + elif [ -x /usr/lib/jvm/java-21-amazon-corretto/bin/java ]; then + export JAVA_HOME=/usr/lib/jvm/java-21-amazon-corretto + elif [ -z "${JAVA_HOME:-}" ]; then + export JAVA_HOME="$(dirname "$(dirname "$(readlink -f "$(command -v java)")")")" + fi + if [ ! -x "$JAVA_HOME/bin/java" ] || [ ! -x "$JAVA_HOME/bin/javac" ]; then + echo "Selected JAVA_HOME=$JAVA_HOME is not a complete JDK" exit 1 fi - export JAVA_HOME="$JAVA_HOME_21_X64" export PATH="$JAVA_HOME/bin:$PATH" + java_spec=$(java -XshowSettings:properties -version 2>&1 | awk '$1 == "java.specification.version" {print $3}') + if [ "$java_spec" != "21" ]; then + echo "Java 21 is required; selected JAVA_HOME=$JAVA_HOME reports specification version $java_spec" + exit 1 + fi + javac_version=$(javac -version 2>&1) + if [[ ! "$javac_version" =~ ^javac[[:space:]]21([.]|[[:space:]]|$) ]]; then + echo "Java 21 javac is required; selected compiler reports $javac_version" + exit 1 + fi java -version + javac -version + if [ -n "${GITHUB_ENV:-}" ]; then + echo "JAVA_HOME=$JAVA_HOME" >> "$GITHUB_ENV" + fi + if [ -n "${GITHUB_PATH:-}" ]; then + echo "$JAVA_HOME/bin" >> "$GITHUB_PATH" + fi prepare_command: | JAVA_SDK_VERSION=$( mvn -B -q \ From c91e61aa53d590f37c06172251140790057a5df0 Mon Sep 17 00:00:00 2001 From: Frank Chen Date: Wed, 7 Oct 2026 22:50:43 +0000 Subject: [PATCH 7/7] ci: pin OTel queue workflow to published squash commit --- .github/workflows/otel-conformance-tests.yml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/.github/workflows/otel-conformance-tests.yml b/.github/workflows/otel-conformance-tests.yml index 452db089d..cfb82bbe7 100644 --- a/.github/workflows/otel-conformance-tests.yml +++ b/.github/workflows/otel-conformance-tests.yml @@ -68,7 +68,7 @@ jobs: actions: write contents: read id-token: write - uses: aws/aws-durable-execution-conformance-tests/.github/workflows/opentelemetry-orchestrator.yml@a66037abbbfa55fde97f714e30f0bc262edefd63 + uses: aws/aws-durable-execution-conformance-tests/.github/workflows/opentelemetry-orchestrator.yml@f18bd0b5f28c5c90e288d0fb8bca08a849b51863 with: runs_on: ${{ github.event_name == 'pull_request' && github.event.pull_request.head.repo.full_name != github.repository && 'ubuntu-latest' || format('codebuild-github-actions-runner-{0}-{1}', github.run_id, github.run_attempt) }} language: java