From 4f5639c25b35c389f90a7af14b147ae251b424e4 Mon Sep 17 00:00:00 2001 From: Alex Wang Date: Mon, 10 Aug 2026 23:38:54 +0000 Subject: [PATCH] fix(plugin): surface missing state fields on hook info records Plugin hook infos dropped state that the Python and JS SDKs carry, so a Java plugin could not see per-operation replay status at the attempt and change hooks, nor operation state at the invocation hooks. Attempt and change-item fields: - UserFunctionStartInfo / UserFunctionEndInfo gain `isReplay`, the operation-level indicator for whether THIS operation was observed via checkpointed state. This is distinct from the existing `isReplayingChildren`, which describes the child operations of a context body and does not substitute for it. - OperationChangeItemInfo was a reduced record; it now carries the full operation surface (`status`, `attempt`, `isReplay`, `error`), ordered to match OperationEndInfo, so an operation seen through a change delta exposes the same fields as through the per-operation hooks. Invocation-info enrichment: - InvocationInfo gains `operations` and `updatedOperations`. - InvocationEndInfo gains `operations` and `executionStartTime`; the latter was present on the start info but dropped from the end record, forcing plugins to correlate back to the start hook. - `updatedOperations` derives from the input's UpdatedOperationIds intersected with the tracked operations, so it is empty on the first invocation and names the externally-completed operations on a replay. - The end-info snapshot is taken at end time, so unlike the start info it also includes operations created during the invocation. ExecutionManager now snapshots the operation ids delivered in the initial state and exposes non-mutating accessors for them. The attempt hook needs the replay indicator at its firing site, but getOperation() delegates to getOperationAndUpdateReplayState, which flips REPLAY to EXECUTION mode as a side effect; reading it from a plugin-hook site would mutate execution state. wasObservedAtInvocationStart is a pure containment check instead, and gives one consistent definition of `isReplay` across the attempt, change and invocation hooks. Both invocation maps are keyed by operation id and valued with OperationChangeItemInfo, now the richest operation snapshot record the SDK has, so a single conversion path feeds the change hook and both invocation hooks. The record name is a wart in this role; renaming it is left as a follow-up. These are positional records, so the added components surface every constructor call site. No defaulting overloads were introduced: a convenience constructor is exactly how a future internal call site would silently ship empty maps to plugins, which is the failure mode this change fixes. Call sites in the OpenTelemetry plugin tests were updated mechanically. Payload surfaces stay out of scope: no `result` and no execution input/result on any info. Verified with the full module build, unit tests and spotless:check. The conformance handlers that assert these field shapes, and the live suite results, are in the stacked follow-up PR. Refs: #604 --- .../durable/otel/ExecutionOtelPluginTest.java | 222 ++++++---- .../otel/InvocationOtelPluginTest.java | 396 +++++++++++------- .../durable/otel/MdcSpanEnricherTest.java | 12 +- .../durable/execution/DurableExecutor.java | 29 +- .../durable/execution/ExecutionManager.java | 54 ++- .../operation/BaseDurableOperation.java | 6 +- .../durable/plugin/InvocationEndInfo.java | 13 + .../lambda/durable/plugin/InvocationInfo.java | 14 +- .../plugin/OperationChangeItemInfo.java | 17 +- .../durable/plugin/PluginInfoConverter.java | 51 ++- .../durable/plugin/UserFunctionEndInfo.java | 4 + .../durable/plugin/UserFunctionStartInfo.java | 4 + .../lambda/durable/DurableConfigTest.java | 4 +- .../plugin/PluginInfoConverterTest.java | 11 +- .../durable/plugin/PluginRunnerTest.java | 19 +- 15 files changed, 593 insertions(+), 263 deletions(-) diff --git a/otel-plugin/src/test/java/software/amazon/lambda/durable/otel/ExecutionOtelPluginTest.java b/otel-plugin/src/test/java/software/amazon/lambda/durable/otel/ExecutionOtelPluginTest.java index fa4fd202e..fa910f718 100644 --- a/otel-plugin/src/test/java/software/amazon/lambda/durable/otel/ExecutionOtelPluginTest.java +++ b/otel-plugin/src/test/java/software/amazon/lambda/durable/otel/ExecutionOtelPluginTest.java @@ -19,6 +19,7 @@ import io.opentelemetry.sdk.trace.SdkTracerProvider; import io.opentelemetry.sdk.trace.export.SimpleSpanProcessor; import java.time.Instant; +import java.util.Map; import java.util.ServiceLoader; import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.BeforeEach; @@ -71,8 +72,9 @@ void customInstrumentationName_isUsedForTracerScope() { .workflowSpanName("Workflow") .instrumentationName("my-custom-scope") .build()); - customPlugin.onInvocationStart(new InvocationInfo("req-1", ARN, true, Instant.now())); - customPlugin.onInvocationEnd(new InvocationEndInfo("req-1", ARN, true, InvocationStatus.SUCCEEDED, null)); + customPlugin.onInvocationStart(new InvocationInfo("req-1", ARN, true, Instant.now(), Map.of(), Map.of())); + customPlugin.onInvocationEnd( + new InvocationEndInfo("req-1", ARN, true, Instant.now(), Map.of(), InvocationStatus.SUCCEEDED, null)); var spans = exporter.getFinishedSpanItems(); assertFalse(spans.isEmpty()); @@ -145,13 +147,14 @@ void defaultConstructor_usesGlobalSdkTracerProviderDirectly() { OpenTelemetrySdk.builder().setTracerProvider(globalTracerProvider).buildAndRegisterGlobal(); var defaultPlugin = new ExecutionOtelPlugin(); - defaultPlugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec1", true, Instant.now())); + defaultPlugin.onInvocationStart( + new InvocationInfo("req-1", "arn:exec1", true, Instant.now(), Map.of(), Map.of())); defaultPlugin.onOperationStart( new OperationInfo("op-1", "step", "STEP", "Step", null, Instant.now(), null, null, false)); defaultPlugin.onOperationEnd(new OperationEndInfo( "op-1", "step", "STEP", "Step", null, Instant.now(), Instant.now(), "SUCCEEDED", null, false, null)); - defaultPlugin.onInvocationEnd( - new InvocationEndInfo("req-1", "arn:exec1", true, InvocationStatus.SUCCEEDED, null)); + defaultPlugin.onInvocationEnd(new InvocationEndInfo( + "req-1", "arn:exec1", true, Instant.now(), Map.of(), InvocationStatus.SUCCEEDED, null)); var spans = globalExporter.getFinishedSpanItems(); // Workflow + Invocation + operation = 3 @@ -178,8 +181,9 @@ void executionOtelPluginProvider_isRegisteredAsServiceProvider() { @Test void terminalInvocation_exportsWorkflowAndInvocationSpans() { - plugin.onInvocationStart(new InvocationInfo("req-1", ARN, true, Instant.now())); - plugin.onInvocationEnd(new InvocationEndInfo("req-1", ARN, true, InvocationStatus.SUCCEEDED, null)); + plugin.onInvocationStart(new InvocationInfo("req-1", ARN, true, Instant.now(), Map.of(), Map.of())); + plugin.onInvocationEnd( + new InvocationEndInfo("req-1", ARN, true, Instant.now(), Map.of(), InvocationStatus.SUCCEEDED, null)); var spans = spanExporter.getFinishedSpanItems(); assertEquals(2, spans.size(), "Terminal invocation should export the Workflow span and the invocation span"); @@ -193,8 +197,9 @@ void terminalInvocation_exportsWorkflowAndInvocationSpans() { @Test void spans_preserveConfiguredServiceName() { - plugin.onInvocationStart(new InvocationInfo("req-1", ARN, true, Instant.now())); - plugin.onInvocationEnd(new InvocationEndInfo("req-1", ARN, true, InvocationStatus.SUCCEEDED, null)); + plugin.onInvocationStart(new InvocationInfo("req-1", ARN, true, Instant.now(), Map.of(), Map.of())); + plugin.onInvocationEnd( + new InvocationEndInfo("req-1", ARN, true, Instant.now(), Map.of(), InvocationStatus.SUCCEEDED, null)); for (var span : spanExporter.getFinishedSpanItems()) { assertEquals( @@ -207,8 +212,9 @@ void spans_preserveConfiguredServiceName() { @Test void workflowSpan_startsAtExecutionStartTime() { var start = Instant.parse("2026-01-15T08:00:00Z"); - plugin.onInvocationStart(new InvocationInfo("req-1", ARN, true, start)); - plugin.onInvocationEnd(new InvocationEndInfo("req-1", ARN, true, InvocationStatus.SUCCEEDED, null)); + plugin.onInvocationStart(new InvocationInfo("req-1", ARN, true, start, Map.of(), Map.of())); + plugin.onInvocationEnd( + new InvocationEndInfo("req-1", ARN, true, Instant.now(), Map.of(), InvocationStatus.SUCCEEDED, null)); var workflowSpan = spanByName(spanExporter.getFinishedSpanItems(), "Workflow"); assertEquals( @@ -219,8 +225,9 @@ void workflowSpan_startsAtExecutionStartTime() { @Test void workflowSpan_hasInternalKind() { - plugin.onInvocationStart(new InvocationInfo("req-1", ARN, true, Instant.now())); - plugin.onInvocationEnd(new InvocationEndInfo("req-1", ARN, true, InvocationStatus.SUCCEEDED, null)); + plugin.onInvocationStart(new InvocationInfo("req-1", ARN, true, Instant.now(), Map.of(), Map.of())); + plugin.onInvocationEnd( + new InvocationEndInfo("req-1", ARN, true, Instant.now(), Map.of(), InvocationStatus.SUCCEEDED, null)); assertEquals( SpanKind.INTERNAL, @@ -230,8 +237,9 @@ void workflowSpan_hasInternalKind() { @Test void workflowAndInvocationSpans_areIndependentRoots_withoutAmbientContext() { - plugin.onInvocationStart(new InvocationInfo("req-1", ARN, true, Instant.now())); - plugin.onInvocationEnd(new InvocationEndInfo("req-1", ARN, true, InvocationStatus.SUCCEEDED, null)); + plugin.onInvocationStart(new InvocationInfo("req-1", ARN, true, Instant.now(), Map.of(), Map.of())); + plugin.onInvocationEnd( + new InvocationEndInfo("req-1", ARN, true, Instant.now(), Map.of(), InvocationStatus.SUCCEEDED, null)); var spans = spanExporter.getFinishedSpanItems(); var workflowSpan = spanByName(spans, "Workflow"); @@ -251,9 +259,10 @@ void invocationStart_usesCurrentSpanContext_whenExtractorReturnsNull() { SpanContext.create(traceId, parentSpanId, TraceFlags.getSampled(), TraceState.getDefault()); try (var ignored = Span.wrap(parentSpanContext).makeCurrent()) { - plugin.onInvocationStart(new InvocationInfo("req-1", ARN, true, Instant.now())); + plugin.onInvocationStart(new InvocationInfo("req-1", ARN, true, Instant.now(), Map.of(), Map.of())); } - plugin.onInvocationEnd(new InvocationEndInfo("req-1", ARN, true, InvocationStatus.SUCCEEDED, null)); + plugin.onInvocationEnd( + new InvocationEndInfo("req-1", ARN, true, Instant.now(), Map.of(), InvocationStatus.SUCCEEDED, null)); var invocationSpan = spanByName(spanExporter.getFinishedSpanItems(), "Invocation"); assertEquals(traceId, invocationSpan.getTraceId()); @@ -262,8 +271,9 @@ void invocationStart_usesCurrentSpanContext_whenExtractorReturnsNull() { @Test void nonTerminalInvocation_doesNotExportWorkflowSpan() { - plugin.onInvocationStart(new InvocationInfo("req-1", ARN, true, Instant.now())); - plugin.onInvocationEnd(new InvocationEndInfo("req-1", ARN, true, InvocationStatus.PENDING, null)); + plugin.onInvocationStart(new InvocationInfo("req-1", ARN, true, Instant.now(), Map.of(), Map.of())); + plugin.onInvocationEnd( + new InvocationEndInfo("req-1", ARN, true, Instant.now(), Map.of(), InvocationStatus.PENDING, null)); var spans = spanExporter.getFinishedSpanItems(); // Only the invocation span is exported; the Workflow span is not ended on non-terminal status. @@ -275,8 +285,9 @@ void nonTerminalInvocation_doesNotExportWorkflowSpan() { @Test void workflowSpan_exportedOnceAcrossInvocations_sameSpanId() { // Invocation 1: non-terminal → no Workflow span exported - plugin.onInvocationStart(new InvocationInfo("req-1", ARN, true, Instant.now())); - plugin.onInvocationEnd(new InvocationEndInfo("req-1", ARN, true, InvocationStatus.PENDING, null)); + plugin.onInvocationStart(new InvocationInfo("req-1", ARN, true, Instant.now(), Map.of(), Map.of())); + plugin.onInvocationEnd( + new InvocationEndInfo("req-1", ARN, true, Instant.now(), Map.of(), InvocationStatus.PENDING, null)); assertTrue( spanExporter.getFinishedSpanItems().stream() .noneMatch(s -> s.getName().equals("Workflow")), @@ -284,8 +295,9 @@ void workflowSpan_exportedOnceAcrossInvocations_sameSpanId() { spanExporter.reset(); // Invocation 2: terminal → Workflow span exported with the deterministic ID - plugin.onInvocationStart(new InvocationInfo("req-2", ARN, false, Instant.now())); - plugin.onInvocationEnd(new InvocationEndInfo("req-2", ARN, false, InvocationStatus.SUCCEEDED, null)); + plugin.onInvocationStart(new InvocationInfo("req-2", ARN, false, Instant.now(), Map.of(), Map.of())); + plugin.onInvocationEnd( + new InvocationEndInfo("req-2", ARN, false, Instant.now(), Map.of(), InvocationStatus.SUCCEEDED, null)); var workflowSpan = spanByName(spanExporter.getFinishedSpanItems(), "Workflow"); assertEquals(StatusCode.OK, workflowSpan.getStatus().getStatusCode()); @@ -296,9 +308,9 @@ void workflowSpan_exportedOnceAcrossInvocations_sameSpanId() { @Test void failedInvocation_setsErrorOnBothWorkflowAndInvocationSpans() { - plugin.onInvocationStart(new InvocationInfo("req-1", ARN, true, Instant.now())); - plugin.onInvocationEnd( - new InvocationEndInfo("req-1", ARN, true, InvocationStatus.FAILED, new RuntimeException("boom"))); + plugin.onInvocationStart(new InvocationInfo("req-1", ARN, true, Instant.now(), Map.of(), Map.of())); + plugin.onInvocationEnd(new InvocationEndInfo( + "req-1", ARN, true, Instant.now(), Map.of(), InvocationStatus.FAILED, new RuntimeException("boom"))); var spans = spanExporter.getFinishedSpanItems(); assertEquals(StatusCode.ERROR, spanByName(spans, "Workflow").getStatus().getStatusCode()); @@ -308,9 +320,15 @@ void failedInvocation_setsErrorOnBothWorkflowAndInvocationSpans() { @Test void retryingInvocation_invocationSpanUnset_workflowNotExported() { - plugin.onInvocationStart(new InvocationInfo("req-1", ARN, true, Instant.now())); + plugin.onInvocationStart(new InvocationInfo("req-1", ARN, true, Instant.now(), Map.of(), Map.of())); plugin.onInvocationEnd(new InvocationEndInfo( - "req-1", ARN, true, InvocationStatus.RETRYING, new RuntimeException("transient"))); + "req-1", + ARN, + true, + Instant.now(), + Map.of(), + InvocationStatus.RETRYING, + new RuntimeException("transient"))); var spans = spanExporter.getFinishedSpanItems(); assertEquals(1, spans.size(), "RETRYING is non-terminal — Workflow span not exported"); @@ -326,12 +344,13 @@ void retryingInvocation_invocationSpanUnset_workflowNotExported() { @Test void operationSpan_carriesAttemptNumberAtEnd() { - plugin.onInvocationStart(new InvocationInfo("req-1", ARN, true, Instant.now())); + plugin.onInvocationStart(new InvocationInfo("req-1", ARN, true, Instant.now(), Map.of(), Map.of())); plugin.onOperationStart( new OperationInfo("op-1", "flaky", "STEP", "Step", null, Instant.now(), null, null, false)); plugin.onOperationEnd(new OperationEndInfo( "op-1", "flaky", "STEP", "Step", null, Instant.now(), Instant.now(), "SUCCEEDED", 3, false, null)); - plugin.onInvocationEnd(new InvocationEndInfo("req-1", ARN, true, InvocationStatus.SUCCEEDED, null)); + plugin.onInvocationEnd( + new InvocationEndInfo("req-1", ARN, true, Instant.now(), Map.of(), InvocationStatus.SUCCEEDED, null)); var operationSpan = spanByName(spanExporter.getFinishedSpanItems(), "flaky"); assertEquals( @@ -344,11 +363,12 @@ void operationSpan_carriesAttemptNumberAtEnd() { @Test void continuationOperationSpan_carriesAttemptNumber() { - plugin.onInvocationStart(new InvocationInfo("req-2", ARN, false, Instant.now())); + plugin.onInvocationStart(new InvocationInfo("req-2", ARN, false, Instant.now(), Map.of(), Map.of())); // No matching onOperationStart in this invocation — continuation branch. plugin.onOperationEnd(new OperationEndInfo( "op-1", "flaky", "STEP", "Step", null, Instant.now(), Instant.now(), "SUCCEEDED", 2, false, null)); - plugin.onInvocationEnd(new InvocationEndInfo("req-2", ARN, false, InvocationStatus.SUCCEEDED, null)); + plugin.onInvocationEnd( + new InvocationEndInfo("req-2", ARN, false, Instant.now(), Map.of(), InvocationStatus.SUCCEEDED, null)); var operationSpan = spanByName(spanExporter.getFinishedSpanItems(), "flaky"); assertEquals( @@ -363,11 +383,12 @@ void continuationOperationSpan_carriesAttemptNumber() { void operationSpan_startsAtOperationStartTimestamp() { var opStart = Instant.parse("2026-02-01T10:00:00Z"); var opEnd = Instant.parse("2026-02-01T10:00:03Z"); - plugin.onInvocationStart(new InvocationInfo("req-1", ARN, true, Instant.now())); + plugin.onInvocationStart(new InvocationInfo("req-1", ARN, true, Instant.now(), Map.of(), Map.of())); plugin.onOperationStart(new OperationInfo("op-1", "step-a", "STEP", "Step", null, opStart, null, null, false)); plugin.onOperationEnd(new OperationEndInfo( "op-1", "step-a", "STEP", "Step", null, opStart, opEnd, "SUCCEEDED", null, false, null)); - plugin.onInvocationEnd(new InvocationEndInfo("req-1", ARN, true, InvocationStatus.SUCCEEDED, null)); + plugin.onInvocationEnd( + new InvocationEndInfo("req-1", ARN, true, Instant.now(), Map.of(), InvocationStatus.SUCCEEDED, null)); var operationSpan = spanByName(spanExporter.getFinishedSpanItems(), "step-a"); assertEquals( @@ -378,12 +399,13 @@ void operationSpan_startsAtOperationStartTimestamp() { @Test void operationSpan_parentedToWorkflow_linkedToInvocation() { - plugin.onInvocationStart(new InvocationInfo("req-1", ARN, true, Instant.now())); + plugin.onInvocationStart(new InvocationInfo("req-1", ARN, true, Instant.now(), Map.of(), Map.of())); plugin.onOperationStart( new OperationInfo("op-1", "step-a", "STEP", "Step", null, Instant.now(), null, null, false)); plugin.onOperationEnd(new OperationEndInfo( "op-1", "step-a", "STEP", "Step", null, Instant.now(), Instant.now(), "SUCCEEDED", null, false, null)); - plugin.onInvocationEnd(new InvocationEndInfo("req-1", ARN, true, InvocationStatus.SUCCEEDED, null)); + plugin.onInvocationEnd( + new InvocationEndInfo("req-1", ARN, true, Instant.now(), Map.of(), InvocationStatus.SUCCEEDED, null)); var spans = spanExporter.getFinishedSpanItems(); var workflowSpan = spanByName(spans, "Workflow"); @@ -402,16 +424,17 @@ void operationSpan_parentedToWorkflow_linkedToInvocation() { @Test void attemptSpan_childOfOperation_linkedToInvocation() { - plugin.onInvocationStart(new InvocationInfo("req-1", ARN, true, Instant.now())); + plugin.onInvocationStart(new InvocationInfo("req-1", ARN, true, Instant.now(), Map.of(), Map.of())); plugin.onOperationStart( new OperationInfo("op-1", "compute", "STEP", "Step", null, Instant.now(), null, null, false)); plugin.onUserFunctionStart( - new UserFunctionStartInfo("op-1", "compute", "STEP", "Step", null, Instant.now(), false, 1)); + new UserFunctionStartInfo("op-1", "compute", "STEP", "Step", null, Instant.now(), false, false, 1)); plugin.onUserFunctionEnd(new UserFunctionEndInfo( - "op-1", "compute", "STEP", "Step", null, Instant.now(), Instant.now(), false, 1, true, null)); + "op-1", "compute", "STEP", "Step", null, Instant.now(), Instant.now(), false, false, 1, true, null)); plugin.onOperationEnd(new OperationEndInfo( "op-1", "compute", "STEP", "Step", null, Instant.now(), Instant.now(), "SUCCEEDED", null, false, null)); - plugin.onInvocationEnd(new InvocationEndInfo("req-1", ARN, true, InvocationStatus.SUCCEEDED, null)); + plugin.onInvocationEnd( + new InvocationEndInfo("req-1", ARN, true, Instant.now(), Map.of(), InvocationStatus.SUCCEEDED, null)); var spans = spanExporter.getFinishedSpanItems(); var operationSpan = spanByName(spans, "compute"); @@ -434,12 +457,24 @@ void attemptSpan_childOfOperation_linkedToInvocation() { @Test void attemptSpan_carriesOperationSubtype() { - plugin.onInvocationStart(new InvocationInfo("req-1", ARN, true, Instant.now())); - plugin.onUserFunctionStart( - new UserFunctionStartInfo("op-1", "process-order", "STEP", "Step", null, Instant.now(), false, 1)); + plugin.onInvocationStart(new InvocationInfo("req-1", ARN, true, Instant.now(), Map.of(), Map.of())); + plugin.onUserFunctionStart(new UserFunctionStartInfo( + "op-1", "process-order", "STEP", "Step", null, Instant.now(), false, false, 1)); plugin.onUserFunctionEnd(new UserFunctionEndInfo( - "op-1", "process-order", "STEP", "Step", null, Instant.now(), Instant.now(), false, 1, true, null)); - plugin.onInvocationEnd(new InvocationEndInfo("req-1", ARN, true, InvocationStatus.SUCCEEDED, null)); + "op-1", + "process-order", + "STEP", + "Step", + null, + Instant.now(), + Instant.now(), + false, + false, + 1, + true, + null)); + plugin.onInvocationEnd( + new InvocationEndInfo("req-1", ARN, true, Instant.now(), Map.of(), InvocationStatus.SUCCEEDED, null)); var attemptSpan = spanExporter.getFinishedSpanItems().stream() .filter(s -> s.getName().contains("process-order")) @@ -453,7 +488,7 @@ void attemptSpan_carriesOperationSubtype() { @Test void childOperation_parentedToParentOperationSpan() { - plugin.onInvocationStart(new InvocationInfo("req-1", ARN, true, Instant.now())); + plugin.onInvocationStart(new InvocationInfo("req-1", ARN, true, Instant.now(), Map.of(), Map.of())); plugin.onOperationStart(new OperationInfo( "op-parent", "my-context", "CONTEXT", "RunInChildContext", null, Instant.now(), null, null, false)); plugin.onOperationStart(new OperationInfo( @@ -482,7 +517,8 @@ void childOperation_parentedToParentOperationSpan() { null, false, null)); - plugin.onInvocationEnd(new InvocationEndInfo("req-1", ARN, true, InvocationStatus.SUCCEEDED, null)); + plugin.onInvocationEnd( + new InvocationEndInfo("req-1", ARN, true, Instant.now(), Map.of(), InvocationStatus.SUCCEEDED, null)); var spans = spanExporter.getFinishedSpanItems(); var parentSpan = spanByName(spans, "my-context"); @@ -497,9 +533,9 @@ void childOperation_parentedToParentOperationSpan() { @Test void userFunctionFailure_setsErrorOnAttemptSpan() { - plugin.onInvocationStart(new InvocationInfo("req-1", ARN, true, Instant.now())); + plugin.onInvocationStart(new InvocationInfo("req-1", ARN, true, Instant.now(), Map.of(), Map.of())); plugin.onUserFunctionStart( - new UserFunctionStartInfo("op-1", "failing", "STEP", "Step", null, Instant.now(), false, 1)); + new UserFunctionStartInfo("op-1", "failing", "STEP", "Step", null, Instant.now(), false, false, 1)); plugin.onUserFunctionEnd(new UserFunctionEndInfo( "op-1", "failing", @@ -509,10 +545,12 @@ void userFunctionFailure_setsErrorOnAttemptSpan() { Instant.now(), Instant.now(), false, + false, 1, false, new RuntimeException("step failed"))); - plugin.onInvocationEnd(new InvocationEndInfo("req-1", ARN, true, InvocationStatus.FAILED, null)); + plugin.onInvocationEnd( + new InvocationEndInfo("req-1", ARN, true, Instant.now(), Map.of(), InvocationStatus.FAILED, null)); var attemptSpan = spanExporter.getFinishedSpanItems().stream() .filter(s -> s.getName().contains("attempt")) @@ -523,12 +561,13 @@ void userFunctionFailure_setsErrorOnAttemptSpan() { @Test void userFunctionSuccess_setsOkOnAttemptSpan() { - plugin.onInvocationStart(new InvocationInfo("req-1", ARN, true, Instant.now())); + plugin.onInvocationStart(new InvocationInfo("req-1", ARN, true, Instant.now(), Map.of(), Map.of())); plugin.onUserFunctionStart( - new UserFunctionStartInfo("op-1", "compute", "STEP", "Step", null, Instant.now(), false, 1)); + new UserFunctionStartInfo("op-1", "compute", "STEP", "Step", null, Instant.now(), false, false, 1)); plugin.onUserFunctionEnd(new UserFunctionEndInfo( - "op-1", "compute", "STEP", "Step", null, Instant.now(), Instant.now(), false, 1, true, null)); - plugin.onInvocationEnd(new InvocationEndInfo("req-1", ARN, true, InvocationStatus.SUCCEEDED, null)); + "op-1", "compute", "STEP", "Step", null, Instant.now(), Instant.now(), false, false, 1, true, null)); + plugin.onInvocationEnd( + new InvocationEndInfo("req-1", ARN, true, Instant.now(), Map.of(), InvocationStatus.SUCCEEDED, null)); var attemptSpan = spanExporter.getFinishedSpanItems().stream() .filter(s -> s.getName().contains("attempt")) @@ -539,12 +578,13 @@ void userFunctionSuccess_setsOkOnAttemptSpan() { @Test void operationSuccess_setsOkOnOperationSpan() { - plugin.onInvocationStart(new InvocationInfo("req-1", ARN, true, Instant.now())); + plugin.onInvocationStart(new InvocationInfo("req-1", ARN, true, Instant.now(), Map.of(), Map.of())); plugin.onOperationStart( new OperationInfo("op-1", "step-ok", "STEP", "Step", null, Instant.now(), null, null, false)); plugin.onOperationEnd(new OperationEndInfo( "op-1", "step-ok", "STEP", "Step", null, Instant.now(), Instant.now(), "SUCCEEDED", null, false, null)); - plugin.onInvocationEnd(new InvocationEndInfo("req-1", ARN, true, InvocationStatus.SUCCEEDED, null)); + plugin.onInvocationEnd( + new InvocationEndInfo("req-1", ARN, true, Instant.now(), Map.of(), InvocationStatus.SUCCEEDED, null)); var operationSpan = spanByName(spanExporter.getFinishedSpanItems(), "step-ok"); assertEquals(StatusCode.OK, operationSpan.getStatus().getStatusCode()); @@ -555,7 +595,7 @@ void operationEnd_withNonSuccessStatusAndNoError_leavesOperationSpanUnset() { // onOperationEnd fires for every terminal status. A CANCELLED operation (or an error-less // FAILED/TIMED_OUT/STOPPED) carries a non-null, non-SUCCEEDED status with a null error. It must NOT be // stamped OK — the span status stays UNSET. - plugin.onInvocationStart(new InvocationInfo("req-1", ARN, true, Instant.now())); + plugin.onInvocationStart(new InvocationInfo("req-1", ARN, true, Instant.now(), Map.of(), Map.of())); plugin.onOperationStart( new OperationInfo("op-cancel", "step-cancel", "STEP", "Step", null, Instant.now(), null, null, false)); plugin.onOperationEnd(new OperationEndInfo( @@ -570,7 +610,8 @@ void operationEnd_withNonSuccessStatusAndNoError_leavesOperationSpanUnset() { null, false, null)); - plugin.onInvocationEnd(new InvocationEndInfo("req-1", ARN, true, InvocationStatus.SUCCEEDED, null)); + plugin.onInvocationEnd( + new InvocationEndInfo("req-1", ARN, true, Instant.now(), Map.of(), InvocationStatus.SUCCEEDED, null)); var operationSpan = spanByName(spanExporter.getFinishedSpanItems(), "step-cancel"); assertEquals(StatusCode.UNSET, operationSpan.getStatus().getStatusCode()); @@ -580,7 +621,7 @@ void operationEnd_withNonSuccessStatusAndNoError_leavesOperationSpanUnset() { void operationEnd_withoutStart_nonSuccessStatusAndNoError_leavesContinuationSpanUnset() { // Same guard on the continuation-span branch (operation completed between invocations): an error-less // TIMED_OUT terminal status must NOT be stamped OK. - plugin.onInvocationStart(new InvocationInfo("req-2", ARN, false, Instant.now())); + plugin.onInvocationStart(new InvocationInfo("req-2", ARN, false, Instant.now(), Map.of(), Map.of())); plugin.onOperationEnd(new OperationEndInfo( "op-cb-timeout", "my-callback", @@ -593,7 +634,8 @@ void operationEnd_withoutStart_nonSuccessStatusAndNoError_leavesContinuationSpan null, false, null)); - plugin.onInvocationEnd(new InvocationEndInfo("req-2", ARN, false, InvocationStatus.SUCCEEDED, null)); + plugin.onInvocationEnd( + new InvocationEndInfo("req-2", ARN, false, Instant.now(), Map.of(), InvocationStatus.SUCCEEDED, null)); var continuationSpan = spanByName(spanExporter.getFinishedSpanItems(), "my-callback"); assertEquals(StatusCode.UNSET, continuationSpan.getStatus().getStatusCode()); @@ -603,12 +645,13 @@ void operationEnd_withoutStart_nonSuccessStatusAndNoError_leavesContinuationSpan void operationEnd_withNullStatusAndNoError_setsOkOnOperationSpan() { // A successful statusless virtual (FLAT CONTEXT) operation fires onOperationEnd with a null operation -> // null status and null error. This is genuine success and must be stamped OK. - plugin.onInvocationStart(new InvocationInfo("req-1", ARN, true, Instant.now())); + plugin.onInvocationStart(new InvocationInfo("req-1", ARN, true, Instant.now(), Map.of(), Map.of())); plugin.onOperationStart( new OperationInfo("op-ctx", "my-ctx", "CONTEXT", null, null, Instant.now(), null, null, false)); plugin.onOperationEnd(new OperationEndInfo( "op-ctx", "my-ctx", "CONTEXT", null, null, Instant.now(), Instant.now(), null, null, false, null)); - plugin.onInvocationEnd(new InvocationEndInfo("req-1", ARN, true, InvocationStatus.SUCCEEDED, null)); + plugin.onInvocationEnd( + new InvocationEndInfo("req-1", ARN, true, Instant.now(), Map.of(), InvocationStatus.SUCCEEDED, null)); var operationSpan = spanByName(spanExporter.getFinishedSpanItems(), "my-ctx"); assertEquals(StatusCode.OK, operationSpan.getStatus().getStatusCode()); @@ -616,10 +659,11 @@ void operationEnd_withNullStatusAndNoError_setsOkOnOperationSpan() { @Test void operationNotCompleted_notEndedAtInvocationEnd() { - plugin.onInvocationStart(new InvocationInfo("req-1", ARN, true, Instant.now())); + plugin.onInvocationStart(new InvocationInfo("req-1", ARN, true, Instant.now(), Map.of(), Map.of())); plugin.onOperationStart( new OperationInfo("op-1", "my-wait", "WAIT", "Wait", null, Instant.now(), null, null, false)); - plugin.onInvocationEnd(new InvocationEndInfo("req-1", ARN, true, InvocationStatus.PENDING, null)); + plugin.onInvocationEnd( + new InvocationEndInfo("req-1", ARN, true, Instant.now(), Map.of(), InvocationStatus.PENDING, null)); var spans = spanExporter.getFinishedSpanItems(); // Only the invocation span is exported. The still-open operation span is NOT force-ended (no PENDING @@ -634,10 +678,11 @@ void operationNotCompleted_notEndedAtInvocationEnd() { @Test void operationOpenedThenCompletedNextInvocation_exportedOnceOnOperationEnd() { // Invocation 1: operation opens but does not complete. - plugin.onInvocationStart(new InvocationInfo("req-1", ARN, true, Instant.now())); + plugin.onInvocationStart(new InvocationInfo("req-1", ARN, true, Instant.now(), Map.of(), Map.of())); plugin.onOperationStart( new OperationInfo("op-1", "my-wait", "WAIT", "Wait", null, Instant.now(), null, null, false)); - plugin.onInvocationEnd(new InvocationEndInfo("req-1", ARN, true, InvocationStatus.PENDING, null)); + plugin.onInvocationEnd( + new InvocationEndInfo("req-1", ARN, true, Instant.now(), Map.of(), InvocationStatus.PENDING, null)); assertTrue( spanExporter.getFinishedSpanItems().stream() .noneMatch(s -> s.getName().equals("my-wait")), @@ -645,10 +690,11 @@ void operationOpenedThenCompletedNextInvocation_exportedOnceOnOperationEnd() { spanExporter.reset(); // Invocation 2: the operation completes → materialized once via onOperationEnd, linked to this invocation. - plugin.onInvocationStart(new InvocationInfo("req-2", ARN, false, Instant.now())); + plugin.onInvocationStart(new InvocationInfo("req-2", ARN, false, Instant.now(), Map.of(), Map.of())); plugin.onOperationEnd(new OperationEndInfo( "op-1", "my-wait", "WAIT", "Wait", null, Instant.now(), Instant.now(), "SUCCEEDED", null, false, null)); - plugin.onInvocationEnd(new InvocationEndInfo("req-2", ARN, false, InvocationStatus.SUCCEEDED, null)); + plugin.onInvocationEnd( + new InvocationEndInfo("req-2", ARN, false, Instant.now(), Map.of(), InvocationStatus.SUCCEEDED, null)); var spans = spanExporter.getFinishedSpanItems(); var waitSpans = @@ -665,17 +711,19 @@ void operationOpenedThenCompletedNextInvocation_exportedOnceOnOperationEnd() { @Test void allSpansShareTraceId_acrossInvocations() { - plugin.onInvocationStart(new InvocationInfo("req-1", ARN, true, Instant.now())); + plugin.onInvocationStart(new InvocationInfo("req-1", ARN, true, Instant.now(), Map.of(), Map.of())); plugin.onOperationStart( new OperationInfo("op-1", "step-1", "STEP", "Step", null, Instant.now(), null, null, false)); plugin.onOperationEnd(new OperationEndInfo( "op-1", "step-1", "STEP", "Step", null, Instant.now(), Instant.now(), "SUCCEEDED", null, false, null)); - plugin.onInvocationEnd(new InvocationEndInfo("req-1", ARN, true, InvocationStatus.PENDING, null)); + plugin.onInvocationEnd( + new InvocationEndInfo("req-1", ARN, true, Instant.now(), Map.of(), InvocationStatus.PENDING, null)); var firstTraceId = spanExporter.getFinishedSpanItems().get(0).getTraceId(); spanExporter.reset(); - plugin.onInvocationStart(new InvocationInfo("req-2", ARN, false, Instant.now())); - plugin.onInvocationEnd(new InvocationEndInfo("req-2", ARN, false, InvocationStatus.SUCCEEDED, null)); + plugin.onInvocationStart(new InvocationInfo("req-2", ARN, false, Instant.now(), Map.of(), Map.of())); + plugin.onInvocationEnd( + new InvocationEndInfo("req-2", ARN, false, Instant.now(), Map.of(), InvocationStatus.SUCCEEDED, null)); var secondSpans = spanExporter.getFinishedSpanItems(); assertTrue( @@ -685,7 +733,7 @@ void allSpansShareTraceId_acrossInvocations() { @Test void operationEnd_withoutStart_createsContinuationSpanWithLink() { - plugin.onInvocationStart(new InvocationInfo("req-2", ARN, false, Instant.now())); + plugin.onInvocationStart(new InvocationInfo("req-2", ARN, false, Instant.now(), Map.of(), Map.of())); // Operation completed between invocations — no matching onOperationStart in this invocation. plugin.onOperationEnd(new OperationEndInfo( "op-wait-1", @@ -699,7 +747,8 @@ void operationEnd_withoutStart_createsContinuationSpanWithLink() { null, false, null)); - plugin.onInvocationEnd(new InvocationEndInfo("req-2", ARN, false, InvocationStatus.SUCCEEDED, null)); + plugin.onInvocationEnd( + new InvocationEndInfo("req-2", ARN, false, Instant.now(), Map.of(), InvocationStatus.SUCCEEDED, null)); var spans = spanExporter.getFinishedSpanItems(); var continuationSpan = spanByName(spans, "my-wait"); @@ -713,8 +762,9 @@ void operationEnd_withoutStart_createsContinuationSpanWithLink() { @Test void deterministicWorkflowSpanId_stableAcrossInvocations() { - plugin.onInvocationStart(new InvocationInfo("req-1", ARN, true, Instant.now())); - plugin.onInvocationEnd(new InvocationEndInfo("req-1", ARN, true, InvocationStatus.SUCCEEDED, null)); + plugin.onInvocationStart(new InvocationInfo("req-1", ARN, true, Instant.now(), Map.of(), Map.of())); + plugin.onInvocationEnd( + new InvocationEndInfo("req-1", ARN, true, Instant.now(), Map.of(), InvocationStatus.SUCCEEDED, null)); var firstWorkflowSpanId = spanByName(spanExporter.getFinishedSpanItems(), "Workflow").getSpanId(); spanExporter.reset(); @@ -728,8 +778,9 @@ void deterministicWorkflowSpanId_stableAcrossInvocations() { .enableMdc(false) .workflowSpanName("Workflow") .build()); - plugin2.onInvocationStart(new InvocationInfo("req-9", ARN, true, Instant.now())); - plugin2.onInvocationEnd(new InvocationEndInfo("req-9", ARN, true, InvocationStatus.SUCCEEDED, null)); + plugin2.onInvocationStart(new InvocationInfo("req-9", ARN, true, Instant.now(), Map.of(), Map.of())); + plugin2.onInvocationEnd( + new InvocationEndInfo("req-9", ARN, true, Instant.now(), Map.of(), InvocationStatus.SUCCEEDED, null)); var secondWorkflowSpanId = spanByName(exporter2.getFinishedSpanItems(), "Workflow").getSpanId(); @@ -753,8 +804,9 @@ void sampling_disabled_producesNoSpans() { .enableMdc(false) .workflowSpanName("Workflow") .build()); - sampledPlugin.onInvocationStart(new InvocationInfo("req-1", ARN, true, Instant.now())); - sampledPlugin.onInvocationEnd(new InvocationEndInfo("req-1", ARN, true, InvocationStatus.SUCCEEDED, null)); + sampledPlugin.onInvocationStart(new InvocationInfo("req-1", ARN, true, Instant.now(), Map.of(), Map.of())); + sampledPlugin.onInvocationEnd( + new InvocationEndInfo("req-1", ARN, true, Instant.now(), Map.of(), InvocationStatus.SUCCEEDED, null)); assertTrue(exporter.getFinishedSpanItems().isEmpty(), "No spans should be exported with 0% sampling"); } @@ -771,12 +823,13 @@ void xrayExtraction_allSpansShareExtractedTraceId() { .enableMdc(false) .workflowSpanName("Workflow") .build()); - xrayPlugin.onInvocationStart(new InvocationInfo("req-1", ARN, true, Instant.now())); + xrayPlugin.onInvocationStart(new InvocationInfo("req-1", ARN, true, Instant.now(), Map.of(), Map.of())); xrayPlugin.onOperationStart( new OperationInfo("op-1", "step-a", "STEP", "Step", null, Instant.now(), null, null, false)); xrayPlugin.onOperationEnd(new OperationEndInfo( "op-1", "step-a", "STEP", "Step", null, Instant.now(), Instant.now(), "SUCCEEDED", null, false, null)); - xrayPlugin.onInvocationEnd(new InvocationEndInfo("req-1", ARN, true, InvocationStatus.SUCCEEDED, null)); + xrayPlugin.onInvocationEnd( + new InvocationEndInfo("req-1", ARN, true, Instant.now(), Map.of(), InvocationStatus.SUCCEEDED, null)); var spans = exporter.getFinishedSpanItems(); assertTrue(spans.size() >= 3, "Workflow + invocation + operation spans expected"); @@ -798,8 +851,9 @@ void xrayExtraction_withParentSpanId_invocationSpanHasCorrectParent() { .workflowSpanName("Workflow") .build()); - xrayPlugin.onInvocationStart(new InvocationInfo("req-1", ARN, true, Instant.now())); - xrayPlugin.onInvocationEnd(new InvocationEndInfo("req-1", ARN, true, InvocationStatus.SUCCEEDED, null)); + xrayPlugin.onInvocationStart(new InvocationInfo("req-1", ARN, true, Instant.now(), Map.of(), Map.of())); + xrayPlugin.onInvocationEnd( + new InvocationEndInfo("req-1", ARN, true, Instant.now(), Map.of(), InvocationStatus.SUCCEEDED, null)); var spans = exporter.getFinishedSpanItems(); var workflowSpan = spanByName(spans, "Workflow"); diff --git a/otel-plugin/src/test/java/software/amazon/lambda/durable/otel/InvocationOtelPluginTest.java b/otel-plugin/src/test/java/software/amazon/lambda/durable/otel/InvocationOtelPluginTest.java index fed449c40..d648c2ae1 100644 --- a/otel-plugin/src/test/java/software/amazon/lambda/durable/otel/InvocationOtelPluginTest.java +++ b/otel-plugin/src/test/java/software/amazon/lambda/durable/otel/InvocationOtelPluginTest.java @@ -28,6 +28,7 @@ import io.opentelemetry.sdk.trace.SdkTracerProviderBuilder; import io.opentelemetry.sdk.trace.export.SimpleSpanProcessor; import java.time.Instant; +import java.util.Map; import java.util.ServiceLoader; import java.util.concurrent.TimeUnit; import java.util.function.BiFunction; @@ -96,13 +97,14 @@ void defaultConstructor_usesGlobalSdkTracerProviderDirectly() { OpenTelemetrySdk.builder().setTracerProvider(globalTracerProvider).buildAndRegisterGlobal(); var defaultPlugin = new InvocationOtelPlugin(); - defaultPlugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec1", true, Instant.now())); + defaultPlugin.onInvocationStart( + new InvocationInfo("req-1", "arn:exec1", true, Instant.now(), Map.of(), Map.of())); defaultPlugin.onOperationStart( new OperationInfo("op-1", "step", "STEP", "Step", null, Instant.now(), null, null, false)); defaultPlugin.onOperationEnd(new OperationEndInfo( "op-1", "step", "STEP", "Step", null, Instant.now(), Instant.now(), "SUCCEEDED", null, false, null)); - defaultPlugin.onInvocationEnd( - new InvocationEndInfo("req-1", "arn:exec1", true, InvocationStatus.SUCCEEDED, null)); + defaultPlugin.onInvocationEnd(new InvocationEndInfo( + "req-1", "arn:exec1", true, Instant.now(), Map.of(), InvocationStatus.SUCCEEDED, null)); var spans = globalExporter.getFinishedSpanItems(); // Plugin creates Workflow + Invocation + operation spans @@ -134,13 +136,14 @@ public ContextPropagators getPropagators() { }); var defaultPlugin = new InvocationOtelPlugin(); - defaultPlugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec1", true, Instant.now())); + defaultPlugin.onInvocationStart( + new InvocationInfo("req-1", "arn:exec1", true, Instant.now(), Map.of(), Map.of())); defaultPlugin.onOperationStart( new OperationInfo("op-1", "step", "STEP", "Step", null, Instant.now(), null, null, false)); defaultPlugin.onOperationEnd(new OperationEndInfo( "op-1", "step", "STEP", "Step", null, Instant.now(), Instant.now(), "SUCCEEDED", null, false, null)); - defaultPlugin.onInvocationEnd( - new InvocationEndInfo("req-1", "arn:exec1", true, InvocationStatus.SUCCEEDED, null)); + defaultPlugin.onInvocationEnd(new InvocationEndInfo( + "req-1", "arn:exec1", true, Instant.now(), Map.of(), InvocationStatus.SUCCEEDED, null)); var spans = globalExporter.getFinishedSpanItems(); // Plugin creates Workflow + Invocation + operation spans @@ -219,9 +222,10 @@ void invocationStart_usesCurrentSpanContext_whenExtractorReturnsNull() { SpanContext.create(traceId, parentSpanId, TraceFlags.getSampled(), TraceState.getDefault()); try (var ignored = Span.wrap(parentSpanContext).makeCurrent()) { - plugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec1", true, Instant.now())); + plugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec1", true, Instant.now(), Map.of(), Map.of())); } - plugin.onInvocationEnd(new InvocationEndInfo("req-1", "arn:exec1", true, InvocationStatus.SUCCEEDED, null)); + plugin.onInvocationEnd(new InvocationEndInfo( + "req-1", "arn:exec1", true, Instant.now(), Map.of(), InvocationStatus.SUCCEEDED, null)); var invocationSpan = spanExporter.getFinishedSpanItems().stream() .filter(span -> span.getName().equals("Invocation")) @@ -234,11 +238,18 @@ void invocationStart_usesCurrentSpanContext_whenExtractorReturnsNull() { @Test void invocationStart_and_end_createsSpan() { plugin.onInvocationStart(new InvocationInfo( - "req-123", "arn:aws:lambda:us-east-1:123:function:test:$LATEST/durable/exec1", true, Instant.now())); + "req-123", + "arn:aws:lambda:us-east-1:123:function:test:$LATEST/durable/exec1", + true, + Instant.now(), + Map.of(), + Map.of())); plugin.onInvocationEnd(new InvocationEndInfo( "req-123", "arn:aws:lambda:us-east-1:123:function:test:$LATEST/durable/exec1", true, + Instant.now(), + Map.of(), InvocationStatus.SUCCEEDED, null)); @@ -261,9 +272,10 @@ void customInstrumentationName_isUsedForTracerScope() { .workflowSpanName("Workflow") .instrumentationName("my-custom-scope") .build()); - customPlugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec1", true, Instant.now())); - customPlugin.onInvocationEnd( - new InvocationEndInfo("req-1", "arn:exec1", true, InvocationStatus.SUCCEEDED, null)); + customPlugin.onInvocationStart( + new InvocationInfo("req-1", "arn:exec1", true, Instant.now(), Map.of(), Map.of())); + customPlugin.onInvocationEnd(new InvocationEndInfo( + "req-1", "arn:exec1", true, Instant.now(), Map.of(), InvocationStatus.SUCCEEDED, null)); var spans = exporter.getFinishedSpanItems(); assertFalse(spans.isEmpty()); @@ -320,8 +332,9 @@ void configOnlyConstructor_rejectsExplicitProviderSource() { @Test void invocationSpan_hasInternalKind() { - plugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec1", true, Instant.now())); - plugin.onInvocationEnd(new InvocationEndInfo("req-1", "arn:exec1", true, InvocationStatus.SUCCEEDED, null)); + plugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec1", true, Instant.now(), Map.of(), Map.of())); + plugin.onInvocationEnd(new InvocationEndInfo( + "req-1", "arn:exec1", true, Instant.now(), Map.of(), InvocationStatus.SUCCEEDED, null)); var span = spanExporter.getFinishedSpanItems().get(0); assertEquals(SpanKind.INTERNAL, span.getKind(), "Invocation span must be INTERNAL kind"); @@ -329,7 +342,7 @@ void invocationSpan_hasInternalKind() { @Test void operationSpanName_usesOperationName_withoutPrefix() { - plugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec1", true, Instant.now())); + plugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec1", true, Instant.now(), Map.of(), Map.of())); plugin.onOperationStart( new OperationInfo("op-1", "create-greeting", "STEP", "Step", null, Instant.now(), null, null, false)); plugin.onOperationEnd(new OperationEndInfo( @@ -344,7 +357,8 @@ void operationSpanName_usesOperationName_withoutPrefix() { null, false, null)); - plugin.onInvocationEnd(new InvocationEndInfo("req-1", "arn:exec1", true, InvocationStatus.SUCCEEDED, null)); + plugin.onInvocationEnd(new InvocationEndInfo( + "req-1", "arn:exec1", true, Instant.now(), Map.of(), InvocationStatus.SUCCEEDED, null)); var operationSpan = spanExporter.getFinishedSpanItems().stream() .filter(s -> s.getName().equals("create-greeting")) @@ -358,12 +372,24 @@ void operationSpanName_usesOperationName_withoutPrefix() { @Test void attemptSpanName_usesOperationNameWithAttemptNumber() { - plugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec1", true, Instant.now())); - plugin.onUserFunctionStart( - new UserFunctionStartInfo("op-1", "process-order", "STEP", "Step", null, Instant.now(), false, 1)); + plugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec1", true, Instant.now(), Map.of(), Map.of())); + plugin.onUserFunctionStart(new UserFunctionStartInfo( + "op-1", "process-order", "STEP", "Step", null, Instant.now(), false, false, 1)); plugin.onUserFunctionEnd(new UserFunctionEndInfo( - "op-1", "process-order", "STEP", "Step", null, Instant.now(), Instant.now(), false, 1, true, null)); - plugin.onInvocationEnd(new InvocationEndInfo("req-1", "arn:exec1", true, InvocationStatus.SUCCEEDED, null)); + "op-1", + "process-order", + "STEP", + "Step", + null, + Instant.now(), + Instant.now(), + false, + false, + 1, + true, + null)); + plugin.onInvocationEnd(new InvocationEndInfo( + "req-1", "arn:exec1", true, Instant.now(), Map.of(), InvocationStatus.SUCCEEDED, null)); var attemptSpan = spanExporter.getFinishedSpanItems().stream() .filter(s -> s.getName().contains("attempt")) @@ -377,13 +403,14 @@ void attemptSpanName_usesOperationNameWithAttemptNumber() { @Test void operationEnd_withAttempt_stampsAttemptNumberOnOperationSpan() { - plugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec1", true, Instant.now())); + plugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec1", true, Instant.now(), Map.of(), Map.of())); plugin.onOperationStart( new OperationInfo("op-1", "flaky", "STEP", "Step", null, Instant.now(), null, null, false)); plugin.onOperationEnd(new OperationEndInfo( "op-1", "flaky", "STEP", "Step", null, Instant.now(), Instant.now(), "SUCCEEDED", 3, false, null)); - plugin.onInvocationEnd(new InvocationEndInfo("req-1", "arn:exec1", true, InvocationStatus.SUCCEEDED, null)); + plugin.onInvocationEnd(new InvocationEndInfo( + "req-1", "arn:exec1", true, Instant.now(), Map.of(), InvocationStatus.SUCCEEDED, null)); var operationSpan = spanExporter.getFinishedSpanItems().stream() .filter(s -> s.getName().equals("flaky")) @@ -397,12 +424,24 @@ void operationEnd_withAttempt_stampsAttemptNumberOnOperationSpan() { @Test void attemptSpan_carriesOperationSubtype() { - plugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec1", true, Instant.now())); - plugin.onUserFunctionStart( - new UserFunctionStartInfo("op-1", "process-order", "STEP", "Step", null, Instant.now(), false, 1)); + plugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec1", true, Instant.now(), Map.of(), Map.of())); + plugin.onUserFunctionStart(new UserFunctionStartInfo( + "op-1", "process-order", "STEP", "Step", null, Instant.now(), false, false, 1)); plugin.onUserFunctionEnd(new UserFunctionEndInfo( - "op-1", "process-order", "STEP", "Step", null, Instant.now(), Instant.now(), false, 1, true, null)); - plugin.onInvocationEnd(new InvocationEndInfo("req-1", "arn:exec1", true, InvocationStatus.SUCCEEDED, null)); + "op-1", + "process-order", + "STEP", + "Step", + null, + Instant.now(), + Instant.now(), + false, + false, + 1, + true, + null)); + plugin.onInvocationEnd(new InvocationEndInfo( + "req-1", "arn:exec1", true, Instant.now(), Map.of(), InvocationStatus.SUCCEEDED, null)); var attemptSpan = spanExporter.getFinishedSpanItems().stream() .filter(s -> s.getName().contains("process-order")) @@ -416,12 +455,13 @@ void attemptSpan_carriesOperationSubtype() { @Test void operationEnd_withoutMatchingStart_stampsAttemptNumberOnContinuationSpan() { - plugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec1", true, Instant.now())); + plugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec1", true, Instant.now(), Map.of(), Map.of())); // No onOperationStart in this invocation → onOperationEnd takes the continuation-span branch. plugin.onOperationEnd(new OperationEndInfo( "op-1", "flaky", "STEP", "Step", null, Instant.now(), Instant.now(), "SUCCEEDED", 2, false, null)); - plugin.onInvocationEnd(new InvocationEndInfo("req-1", "arn:exec1", true, InvocationStatus.SUCCEEDED, null)); + plugin.onInvocationEnd(new InvocationEndInfo( + "req-1", "arn:exec1", true, Instant.now(), Map.of(), InvocationStatus.SUCCEEDED, null)); var continuationSpan = spanExporter.getFinishedSpanItems().stream() .filter(s -> s.getName().equals("flaky")) @@ -439,9 +479,15 @@ void operationEnd_withoutMatchingStart_stampsAttemptNumberOnContinuationSpan() { @Test void invocationEnd_withFailure_setsErrorStatus() { - plugin.onInvocationStart(new InvocationInfo("req-123", "arn:exec1", true, Instant.now())); + plugin.onInvocationStart(new InvocationInfo("req-123", "arn:exec1", true, Instant.now(), Map.of(), Map.of())); plugin.onInvocationEnd(new InvocationEndInfo( - "req-123", "arn:exec1", true, InvocationStatus.FAILED, new RuntimeException("boom"))); + "req-123", + "arn:exec1", + true, + Instant.now(), + Map.of(), + InvocationStatus.FAILED, + new RuntimeException("boom"))); var spans = spanExporter.getFinishedSpanItems(); assertEquals(2, spans.size()); // invocation + Workflow @@ -450,9 +496,15 @@ void invocationEnd_withFailure_setsErrorStatus() { @Test void invocationEnd_withRetrying_leavesStatusUnset() { - plugin.onInvocationStart(new InvocationInfo("req-123", "arn:exec1", true, Instant.now())); + plugin.onInvocationStart(new InvocationInfo("req-123", "arn:exec1", true, Instant.now(), Map.of(), Map.of())); plugin.onInvocationEnd(new InvocationEndInfo( - "req-123", "arn:exec1", true, InvocationStatus.RETRYING, new RuntimeException("transient"))); + "req-123", + "arn:exec1", + true, + Instant.now(), + Map.of(), + InvocationStatus.RETRYING, + new RuntimeException("transient"))); var spans = spanExporter.getFinishedSpanItems(); assertEquals(1, spans.size()); @@ -464,7 +516,7 @@ void invocationEnd_withRetrying_leavesStatusUnset() { @Test void operationStart_createsSpan_operationEnd_endsIt() { - plugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec1", true, Instant.now())); + plugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec1", true, Instant.now(), Map.of(), Map.of())); var start = Instant.parse("2026-06-01T10:00:00Z"); var end = Instant.parse("2026-06-01T10:00:05Z"); @@ -477,7 +529,8 @@ void operationStart_createsSpan_operationEnd_endsIt() { plugin.onOperationEnd(new OperationEndInfo( "op-hash-1", "my-step", "STEP", "Step", null, start, end, "SUCCEEDED", null, false, null)); - plugin.onInvocationEnd(new InvocationEndInfo("req-1", "arn:exec1", true, InvocationStatus.SUCCEEDED, null)); + plugin.onInvocationEnd(new InvocationEndInfo( + "req-1", "arn:exec1", true, Instant.now(), Map.of(), InvocationStatus.SUCCEEDED, null)); var spans = spanExporter.getFinishedSpanItems(); assertEquals(3, spans.size()); // operation + invocation + Workflow @@ -491,15 +544,16 @@ void operationStart_createsSpan_operationEnd_endsIt() { @Test void userFunctionStart_and_end_createsAttemptSpan() { - plugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec1", true, Instant.now())); + plugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec1", true, Instant.now(), Map.of(), Map.of())); plugin.onUserFunctionStart( - new UserFunctionStartInfo("op-1", "compute", "STEP", "Step", null, Instant.now(), false, 1)); + new UserFunctionStartInfo("op-1", "compute", "STEP", "Step", null, Instant.now(), false, false, 1)); plugin.onUserFunctionEnd(new UserFunctionEndInfo( - "op-1", "compute", "STEP", "Step", null, Instant.now(), Instant.now(), false, 1, true, null)); + "op-1", "compute", "STEP", "Step", null, Instant.now(), Instant.now(), false, false, 1, true, null)); - plugin.onInvocationEnd(new InvocationEndInfo("req-1", "arn:exec1", true, InvocationStatus.SUCCEEDED, null)); + plugin.onInvocationEnd(new InvocationEndInfo( + "req-1", "arn:exec1", true, Instant.now(), Map.of(), InvocationStatus.SUCCEEDED, null)); var spans = spanExporter.getFinishedSpanItems(); assertEquals(3, spans.size()); // attempt + invocation + Workflow @@ -514,10 +568,10 @@ void userFunctionStart_and_end_createsAttemptSpan() { @Test void userFunctionEnd_withFailure_setsErrorOnAttemptSpan() { - plugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec1", true, Instant.now())); + plugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec1", true, Instant.now(), Map.of(), Map.of())); plugin.onUserFunctionStart( - new UserFunctionStartInfo("op-1", "failing", "STEP", "Step", null, Instant.now(), false, 1)); + new UserFunctionStartInfo("op-1", "failing", "STEP", "Step", null, Instant.now(), false, false, 1)); plugin.onUserFunctionEnd(new UserFunctionEndInfo( "op-1", @@ -528,11 +582,13 @@ void userFunctionEnd_withFailure_setsErrorOnAttemptSpan() { Instant.now(), Instant.now(), false, + false, 1, false, new RuntimeException("step failed"))); - plugin.onInvocationEnd(new InvocationEndInfo("req-1", "arn:exec1", true, InvocationStatus.FAILED, null)); + plugin.onInvocationEnd(new InvocationEndInfo( + "req-1", "arn:exec1", true, Instant.now(), Map.of(), InvocationStatus.FAILED, null)); var attemptSpan = spanExporter.getFinishedSpanItems().stream() .filter(s -> s.getName().contains("attempt")) @@ -543,14 +599,15 @@ void userFunctionEnd_withFailure_setsErrorOnAttemptSpan() { @Test void userFunctionEnd_withSuccess_setsOkOnAttemptSpan() { - plugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec1", true, Instant.now())); + plugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec1", true, Instant.now(), Map.of(), Map.of())); plugin.onUserFunctionStart( - new UserFunctionStartInfo("op-1", "compute", "STEP", "Step", null, Instant.now(), false, 1)); + new UserFunctionStartInfo("op-1", "compute", "STEP", "Step", null, Instant.now(), false, false, 1)); plugin.onUserFunctionEnd(new UserFunctionEndInfo( - "op-1", "compute", "STEP", "Step", null, Instant.now(), Instant.now(), false, 1, true, null)); + "op-1", "compute", "STEP", "Step", null, Instant.now(), Instant.now(), false, false, 1, true, null)); - plugin.onInvocationEnd(new InvocationEndInfo("req-1", "arn:exec1", true, InvocationStatus.SUCCEEDED, null)); + plugin.onInvocationEnd(new InvocationEndInfo( + "req-1", "arn:exec1", true, Instant.now(), Map.of(), InvocationStatus.SUCCEEDED, null)); var attemptSpan = spanExporter.getFinishedSpanItems().stream() .filter(s -> s.getName().contains("attempt")) @@ -561,14 +618,15 @@ void userFunctionEnd_withSuccess_setsOkOnAttemptSpan() { @Test void operationEnd_withSuccess_setsOkOnOperationSpan() { - plugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec1", true, Instant.now())); + plugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec1", true, Instant.now(), Map.of(), Map.of())); plugin.onOperationStart( new OperationInfo("op-1", "step-ok", "STEP", "Step", null, Instant.now(), null, null, false)); plugin.onOperationEnd(new OperationEndInfo( "op-1", "step-ok", "STEP", "Step", null, Instant.now(), Instant.now(), "SUCCEEDED", null, false, null)); - plugin.onInvocationEnd(new InvocationEndInfo("req-1", "arn:exec1", true, InvocationStatus.SUCCEEDED, null)); + plugin.onInvocationEnd(new InvocationEndInfo( + "req-1", "arn:exec1", true, Instant.now(), Map.of(), InvocationStatus.SUCCEEDED, null)); var operationSpan = spanExporter.getFinishedSpanItems().stream() .filter(s -> "step-ok".equals(s.getName())) @@ -584,7 +642,7 @@ void operationEnd_withNonSuccessStatusAndNoError_leavesOperationSpanUnset() { // onOperationEnd fires for every terminal status. A CANCELLED operation (or an error-less // FAILED/TIMED_OUT/STOPPED) carries a non-null, non-SUCCEEDED status with a null error. It must NOT be // stamped OK — the span status stays UNSET. - plugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec1", true, Instant.now())); + plugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec1", true, Instant.now(), Map.of(), Map.of())); plugin.onOperationStart( new OperationInfo("op-cancel", "step-cancel", "STEP", "Step", null, Instant.now(), null, null, false)); @@ -601,7 +659,8 @@ void operationEnd_withNonSuccessStatusAndNoError_leavesOperationSpanUnset() { false, null)); - plugin.onInvocationEnd(new InvocationEndInfo("req-1", "arn:exec1", true, InvocationStatus.SUCCEEDED, null)); + plugin.onInvocationEnd(new InvocationEndInfo( + "req-1", "arn:exec1", true, Instant.now(), Map.of(), InvocationStatus.SUCCEEDED, null)); var operationSpan = spanExporter.getFinishedSpanItems().stream() .filter(s -> "step-cancel".equals(s.getName())) @@ -614,7 +673,7 @@ void operationEnd_withNonSuccessStatusAndNoError_leavesOperationSpanUnset() { void operationEnd_withoutMatchingStart_nonSuccessStatusAndNoError_leavesContinuationSpanUnset() { // Same guard on the continuation-span branch (operation completed between invocations): an error-less // TIMED_OUT terminal status must NOT be stamped OK. - plugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec1", true, Instant.now())); + plugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec1", true, Instant.now(), Map.of(), Map.of())); plugin.onOperationEnd(new OperationEndInfo( "op-cb-timeout", @@ -629,7 +688,8 @@ void operationEnd_withoutMatchingStart_nonSuccessStatusAndNoError_leavesContinua false, null)); - plugin.onInvocationEnd(new InvocationEndInfo("req-1", "arn:exec1", true, InvocationStatus.SUCCEEDED, null)); + plugin.onInvocationEnd(new InvocationEndInfo( + "req-1", "arn:exec1", true, Instant.now(), Map.of(), InvocationStatus.SUCCEEDED, null)); var continuationSpan = spanExporter.getFinishedSpanItems().stream() .filter(s -> s.getName().contains("callback")) @@ -642,14 +702,15 @@ void operationEnd_withoutMatchingStart_nonSuccessStatusAndNoError_leavesContinua void operationEnd_withNullStatusAndNoError_setsOkOnOperationSpan() { // A successful statusless virtual (FLAT CONTEXT) operation fires onOperationEnd with a null operation -> // null status and null error. This is genuine success and must be stamped OK. - plugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec1", true, Instant.now())); + plugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec1", true, Instant.now(), Map.of(), Map.of())); plugin.onOperationStart( new OperationInfo("op-ctx", "my-ctx", "CONTEXT", null, null, Instant.now(), null, null, false)); plugin.onOperationEnd(new OperationEndInfo( "op-ctx", "my-ctx", "CONTEXT", null, null, Instant.now(), Instant.now(), null, null, false, null)); - plugin.onInvocationEnd(new InvocationEndInfo("req-1", "arn:exec1", true, InvocationStatus.SUCCEEDED, null)); + plugin.onInvocationEnd(new InvocationEndInfo( + "req-1", "arn:exec1", true, Instant.now(), Map.of(), InvocationStatus.SUCCEEDED, null)); var operationSpan = spanExporter.getFinishedSpanItems().stream() .filter(s -> "my-ctx".equals(s.getName())) @@ -661,15 +722,15 @@ void operationEnd_withNullStatusAndNoError_setsOkOnOperationSpan() { @Test void fullLifecycle_producesCorrectSpanHierarchy() { var arn = "arn:aws:lambda:us-east-1:123:function:test:$LATEST/durable/exec1"; - plugin.onInvocationStart(new InvocationInfo("req-1", arn, true, Instant.now())); + plugin.onInvocationStart(new InvocationInfo("req-1", arn, true, Instant.now(), Map.of(), Map.of())); // Step 1: operation starts, user function runs, operation completes plugin.onOperationStart( new OperationInfo("op-1", "step-a", "STEP", "Step", null, Instant.now(), null, null, false)); plugin.onUserFunctionStart( - new UserFunctionStartInfo("op-1", "step-a", "STEP", "Step", null, Instant.now(), false, 1)); + new UserFunctionStartInfo("op-1", "step-a", "STEP", "Step", null, Instant.now(), false, false, 1)); plugin.onUserFunctionEnd(new UserFunctionEndInfo( - "op-1", "step-a", "STEP", "Step", null, Instant.now(), Instant.now(), false, 1, true, null)); + "op-1", "step-a", "STEP", "Step", null, Instant.now(), Instant.now(), false, false, 1, true, null)); plugin.onOperationEnd(new OperationEndInfo( "op-1", "step-a", "STEP", "Step", null, Instant.now(), Instant.now(), "SUCCEEDED", null, false, null)); @@ -677,13 +738,14 @@ void fullLifecycle_producesCorrectSpanHierarchy() { plugin.onOperationStart( new OperationInfo("op-2", "step-b", "STEP", "Step", null, Instant.now(), null, null, false)); plugin.onUserFunctionStart( - new UserFunctionStartInfo("op-2", "step-b", "STEP", "Step", null, Instant.now(), false, 1)); + new UserFunctionStartInfo("op-2", "step-b", "STEP", "Step", null, Instant.now(), false, false, 1)); plugin.onUserFunctionEnd(new UserFunctionEndInfo( - "op-2", "step-b", "STEP", "Step", null, Instant.now(), Instant.now(), false, 1, true, null)); + "op-2", "step-b", "STEP", "Step", null, Instant.now(), Instant.now(), false, false, 1, true, null)); plugin.onOperationEnd(new OperationEndInfo( "op-2", "step-b", "STEP", "Step", null, Instant.now(), Instant.now(), "SUCCEEDED", null, false, null)); - plugin.onInvocationEnd(new InvocationEndInfo("req-1", arn, true, InvocationStatus.SUCCEEDED, null)); + plugin.onInvocationEnd( + new InvocationEndInfo("req-1", arn, true, Instant.now(), Map.of(), InvocationStatus.SUCCEEDED, null)); var spans = spanExporter.getFinishedSpanItems(); // 2 attempt spans + 2 operation spans + 1 invocation span + 1 Workflow span = 6 @@ -698,15 +760,17 @@ void fullLifecycle_producesCorrectSpanHierarchy() { void deterministicIds_sameExecutionProducesSameTraceId() { var arn = "arn:aws:lambda:us-east-1:123:function:test:$LATEST/durable/exec1"; - plugin.onInvocationStart(new InvocationInfo("req-1", arn, true, Instant.now())); - plugin.onInvocationEnd(new InvocationEndInfo("req-1", arn, true, InvocationStatus.PENDING, null)); + plugin.onInvocationStart(new InvocationInfo("req-1", arn, true, Instant.now(), Map.of(), Map.of())); + plugin.onInvocationEnd( + new InvocationEndInfo("req-1", arn, true, Instant.now(), Map.of(), InvocationStatus.PENDING, null)); var firstTraceId = spanExporter.getFinishedSpanItems().get(0).getTraceId(); spanExporter.reset(); // Second invocation of same execution - plugin.onInvocationStart(new InvocationInfo("req-2", arn, false, Instant.now())); - plugin.onInvocationEnd(new InvocationEndInfo("req-2", arn, false, InvocationStatus.SUCCEEDED, null)); + plugin.onInvocationStart(new InvocationInfo("req-2", arn, false, Instant.now(), Map.of(), Map.of())); + plugin.onInvocationEnd( + new InvocationEndInfo("req-2", arn, false, Instant.now(), Map.of(), InvocationStatus.SUCCEEDED, null)); var secondTraceId = spanExporter.getFinishedSpanItems().get(0).getTraceId(); @@ -715,14 +779,15 @@ void deterministicIds_sameExecutionProducesSameTraceId() { @Test void operationNotCompleted_spanEndedAtInvocationEnd() { - plugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec1", true, Instant.now())); + plugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec1", true, Instant.now(), Map.of(), Map.of())); // Operation starts but never completes (e.g., wait operation, invocation suspends) plugin.onOperationStart( new OperationInfo("op-1", "my-wait", "WAIT", "Wait", null, Instant.now(), null, null, false)); // Invocation ends without onOperationEnd being called - plugin.onInvocationEnd(new InvocationEndInfo("req-1", "arn:exec1", true, InvocationStatus.PENDING, null)); + plugin.onInvocationEnd(new InvocationEndInfo( + "req-1", "arn:exec1", true, Instant.now(), Map.of(), InvocationStatus.PENDING, null)); var spans = spanExporter.getFinishedSpanItems(); // Should have: operation span (ended at invocation end) + invocation span @@ -738,11 +803,12 @@ void operationNotCompleted_spanEndedAtInvocationEnd() { @Test void operationStart_withStatus_preservesStatus() { - plugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec1", false, Instant.now())); + plugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec1", false, Instant.now(), Map.of(), Map.of())); plugin.onOperationStart( new OperationInfo("op-1", "my-step", "STEP", "Step", null, Instant.now(), null, "PENDING", true)); - plugin.onInvocationEnd(new InvocationEndInfo("req-1", "arn:exec1", false, InvocationStatus.PENDING, null)); + plugin.onInvocationEnd(new InvocationEndInfo( + "req-1", "arn:exec1", false, Instant.now(), Map.of(), InvocationStatus.PENDING, null)); var operationSpan = spanExporter.getFinishedSpanItems().stream() .filter(s -> "my-step".equals(s.getName())) @@ -755,15 +821,16 @@ void operationStart_withStatus_preservesStatus() { void invocationEnd_closesNestedSpansChildFirst() { var parentId = "op-parent"; var childId = "op-child"; - plugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec1", true, Instant.now())); + plugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec1", true, Instant.now(), Map.of(), Map.of())); plugin.onOperationStart(new OperationInfo( parentId, "parent-context", "CONTEXT", "RunInChildContext", null, Instant.now(), null, null, false)); plugin.onOperationStart( new OperationInfo(childId, "child-step", "STEP", "Step", parentId, Instant.now(), null, null, false)); - plugin.onUserFunctionStart( - new UserFunctionStartInfo(childId, "child-step", "STEP", "Step", parentId, Instant.now(), false, 1)); + plugin.onUserFunctionStart(new UserFunctionStartInfo( + childId, "child-step", "STEP", "Step", parentId, Instant.now(), false, false, 1)); - plugin.onInvocationEnd(new InvocationEndInfo("req-1", "arn:exec1", true, InvocationStatus.PENDING, null)); + plugin.onInvocationEnd(new InvocationEndInfo( + "req-1", "arn:exec1", true, Instant.now(), Map.of(), InvocationStatus.PENDING, null)); var parentSpan = spanByName("parent-context"); var childSpan = spanByName("child-step"); @@ -795,15 +862,16 @@ void sampling_disabled_producesNoSpans() { .enableMdc(false) .build()); - sampledPlugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec1", true, Instant.now())); + sampledPlugin.onInvocationStart( + new InvocationInfo("req-1", "arn:exec1", true, Instant.now(), Map.of(), Map.of())); sampledPlugin.onUserFunctionStart( - new UserFunctionStartInfo("op-1", "step", "STEP", "Step", null, Instant.now(), false, 1)); + new UserFunctionStartInfo("op-1", "step", "STEP", "Step", null, Instant.now(), false, false, 1)); sampledPlugin.onUserFunctionEnd(new UserFunctionEndInfo( - "op-1", "step", "STEP", "Step", null, Instant.now(), Instant.now(), false, 1, true, null)); + "op-1", "step", "STEP", "Step", null, Instant.now(), Instant.now(), false, false, 1, true, null)); sampledPlugin.onOperationEnd(new OperationEndInfo( "op-1", "step", "STEP", "Step", null, Instant.now(), Instant.now(), "SUCCEEDED", null, false, null)); - sampledPlugin.onInvocationEnd( - new InvocationEndInfo("req-1", "arn:exec1", true, InvocationStatus.SUCCEEDED, null)); + sampledPlugin.onInvocationEnd(new InvocationEndInfo( + "req-1", "arn:exec1", true, Instant.now(), Map.of(), InvocationStatus.SUCCEEDED, null)); assertTrue(spanExporter.getFinishedSpanItems().isEmpty(), "No spans should be exported with 0% sampling"); } @@ -823,8 +891,9 @@ void xrayExtraction_usesExtractedTraceId_overArnDerived() { .enableMdc(false) .build()); - xrayPlugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec1", true, Instant.now())); - xrayPlugin.onInvocationEnd(new InvocationEndInfo("req-1", "arn:exec1", true, InvocationStatus.SUCCEEDED, null)); + xrayPlugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec1", true, Instant.now(), Map.of(), Map.of())); + xrayPlugin.onInvocationEnd(new InvocationEndInfo( + "req-1", "arn:exec1", true, Instant.now(), Map.of(), InvocationStatus.SUCCEEDED, null)); var spans = spanExporter.getFinishedSpanItems(); assertEquals(2, spans.size()); // invocation + Workflow @@ -844,16 +913,17 @@ void xrayExtraction_allSpansShareExtractedTraceId() { .enableMdc(false) .build()); - xrayPlugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec1", true, Instant.now())); + xrayPlugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec1", true, Instant.now(), Map.of(), Map.of())); xrayPlugin.onOperationStart( new OperationInfo("op-1", "step-a", "STEP", "Step", null, Instant.now(), null, null, false)); xrayPlugin.onUserFunctionStart( - new UserFunctionStartInfo("op-1", "step-a", "STEP", "Step", null, Instant.now(), false, 1)); + new UserFunctionStartInfo("op-1", "step-a", "STEP", "Step", null, Instant.now(), false, false, 1)); xrayPlugin.onUserFunctionEnd(new UserFunctionEndInfo( - "op-1", "step-a", "STEP", "Step", null, Instant.now(), Instant.now(), false, 1, true, null)); + "op-1", "step-a", "STEP", "Step", null, Instant.now(), Instant.now(), false, false, 1, true, null)); xrayPlugin.onOperationEnd(new OperationEndInfo( "op-1", "step-a", "STEP", "Step", null, Instant.now(), Instant.now(), "SUCCEEDED", null, false, null)); - xrayPlugin.onInvocationEnd(new InvocationEndInfo("req-1", "arn:exec1", true, InvocationStatus.SUCCEEDED, null)); + xrayPlugin.onInvocationEnd(new InvocationEndInfo( + "req-1", "arn:exec1", true, Instant.now(), Map.of(), InvocationStatus.SUCCEEDED, null)); var spans = spanExporter.getFinishedSpanItems(); assertTrue(spans.size() >= 2, "Should have invocation + operation + attempt spans"); @@ -876,8 +946,9 @@ void xrayExtraction_withParentSpanId_invocationSpanHasCorrectParent() { .enableMdc(false) .build()); - xrayPlugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec1", true, Instant.now())); - xrayPlugin.onInvocationEnd(new InvocationEndInfo("req-1", "arn:exec1", true, InvocationStatus.SUCCEEDED, null)); + xrayPlugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec1", true, Instant.now(), Map.of(), Map.of())); + xrayPlugin.onInvocationEnd(new InvocationEndInfo( + "req-1", "arn:exec1", true, Instant.now(), Map.of(), InvocationStatus.SUCCEEDED, null)); var spans = spanExporter.getFinishedSpanItems(); assertEquals(2, spans.size()); // invocation + Workflow @@ -903,8 +974,9 @@ void xrayExtraction_withoutParentSpanId_invocationSpanIsRoot() { .enableMdc(false) .build()); - xrayPlugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec1", true, Instant.now())); - xrayPlugin.onInvocationEnd(new InvocationEndInfo("req-1", "arn:exec1", true, InvocationStatus.SUCCEEDED, null)); + xrayPlugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec1", true, Instant.now(), Map.of(), Map.of())); + xrayPlugin.onInvocationEnd(new InvocationEndInfo( + "req-1", "arn:exec1", true, Instant.now(), Map.of(), InvocationStatus.SUCCEEDED, null)); var spans = spanExporter.getFinishedSpanItems(); assertEquals(2, spans.size()); // invocation + Workflow @@ -936,21 +1008,23 @@ void xrayExtraction_multipleInvocations_sameTraceId_unifiedTrace() { .build()); // First invocation - xrayPlugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec1", true, Instant.now())); + xrayPlugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec1", true, Instant.now(), Map.of(), Map.of())); xrayPlugin.onOperationStart( new OperationInfo("op-1", "step-1", "STEP", "Step", null, Instant.now(), null, null, false)); xrayPlugin.onOperationEnd(new OperationEndInfo( "op-1", "step-1", "STEP", "Step", null, Instant.now(), Instant.now(), "SUCCEEDED", null, false, null)); - xrayPlugin.onInvocationEnd(new InvocationEndInfo("req-1", "arn:exec1", true, InvocationStatus.PENDING, null)); + xrayPlugin.onInvocationEnd(new InvocationEndInfo( + "req-1", "arn:exec1", true, Instant.now(), Map.of(), InvocationStatus.PENDING, null)); // Second invocation (same execution, same X-Ray Root from backend) - xrayPlugin.onInvocationStart(new InvocationInfo("req-2", "arn:exec1", false, Instant.now())); + xrayPlugin.onInvocationStart( + new InvocationInfo("req-2", "arn:exec1", false, Instant.now(), Map.of(), Map.of())); xrayPlugin.onOperationStart( new OperationInfo("op-2", "step-2", "STEP", "Step", null, Instant.now(), null, null, false)); xrayPlugin.onOperationEnd(new OperationEndInfo( "op-2", "step-2", "STEP", "Step", null, Instant.now(), Instant.now(), "SUCCEEDED", null, false, null)); - xrayPlugin.onInvocationEnd( - new InvocationEndInfo("req-2", "arn:exec1", false, InvocationStatus.SUCCEEDED, null)); + xrayPlugin.onInvocationEnd(new InvocationEndInfo( + "req-2", "arn:exec1", false, Instant.now(), Map.of(), InvocationStatus.SUCCEEDED, null)); var spans = spanExporter.getFinishedSpanItems(); assertTrue(spans.size() >= 4, "Should have spans from both invocations"); @@ -972,8 +1046,9 @@ void xrayExtraction_nullExtractor_fallsBackToArnDerived() { .build()); var arn = "arn:aws:lambda:us-east-1:123:function:test:$LATEST/durable/exec1"; - noXrayPlugin.onInvocationStart(new InvocationInfo("req-1", arn, true, Instant.now())); - noXrayPlugin.onInvocationEnd(new InvocationEndInfo("req-1", arn, true, InvocationStatus.SUCCEEDED, null)); + noXrayPlugin.onInvocationStart(new InvocationInfo("req-1", arn, true, Instant.now(), Map.of(), Map.of())); + noXrayPlugin.onInvocationEnd( + new InvocationEndInfo("req-1", arn, true, Instant.now(), Map.of(), InvocationStatus.SUCCEEDED, null)); var spans = spanExporter.getFinishedSpanItems(); assertEquals(2, spans.size()); // invocation + Workflow @@ -1004,8 +1079,9 @@ void xrayExtraction_extractedTraceIdMatchesXrayConversion() { .enableMdc(false) .build()); - xrayPlugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec1", true, Instant.now())); - xrayPlugin.onInvocationEnd(new InvocationEndInfo("req-1", "arn:exec1", true, InvocationStatus.SUCCEEDED, null)); + xrayPlugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec1", true, Instant.now(), Map.of(), Map.of())); + xrayPlugin.onInvocationEnd(new InvocationEndInfo( + "req-1", "arn:exec1", true, Instant.now(), Map.of(), InvocationStatus.SUCCEEDED, null)); var spans = spanExporter.getFinishedSpanItems(); assertEquals(expectedOtelTraceId, spans.get(0).getTraceId()); @@ -1015,7 +1091,7 @@ void xrayExtraction_extractedTraceIdMatchesXrayConversion() { @Test void operationEnd_withoutMatchingStart_createsContinuationSpanWithLink() { - plugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec1", true, Instant.now())); + plugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec1", true, Instant.now(), Map.of(), Map.of())); // onOperationEnd without a prior onOperationStart — operation completed between invocations plugin.onOperationEnd(new OperationEndInfo( @@ -1031,7 +1107,8 @@ void operationEnd_withoutMatchingStart_createsContinuationSpanWithLink() { false, null)); - plugin.onInvocationEnd(new InvocationEndInfo("req-1", "arn:exec1", true, InvocationStatus.SUCCEEDED, null)); + plugin.onInvocationEnd(new InvocationEndInfo( + "req-1", "arn:exec1", true, Instant.now(), Map.of(), InvocationStatus.SUCCEEDED, null)); var spans = spanExporter.getFinishedSpanItems(); assertEquals(3, spans.size()); // continuation + invocation + Workflow @@ -1046,7 +1123,7 @@ void operationEnd_withoutMatchingStart_createsContinuationSpanWithLink() { @Test void operationEnd_withoutMatchingStart_startsWithinCurrentInvocation() { - plugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec1", true, Instant.now())); + plugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec1", true, Instant.now(), Map.of(), Map.of())); var operationStart = Instant.EPOCH; var operationEnd = operationStart.plusSeconds(60); @@ -1065,7 +1142,8 @@ void operationEnd_withoutMatchingStart_startsWithinCurrentInvocation() { false, null)); - plugin.onInvocationEnd(new InvocationEndInfo("req-1", "arn:exec1", true, InvocationStatus.SUCCEEDED, null)); + plugin.onInvocationEnd(new InvocationEndInfo( + "req-1", "arn:exec1", true, Instant.now(), Map.of(), InvocationStatus.SUCCEEDED, null)); var spans = spanExporter.getFinishedSpanItems(); var continuationSpan = spans.stream() @@ -1084,7 +1162,7 @@ void operationEnd_withoutMatchingStart_startsWithinCurrentInvocation() { @Test void operationEnd_withoutMatchingStart_withError_setsErrorStatus() { - plugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec1", true, Instant.now())); + plugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec1", true, Instant.now(), Map.of(), Map.of())); plugin.onOperationEnd(new OperationEndInfo( "op-cb-1", @@ -1099,7 +1177,8 @@ void operationEnd_withoutMatchingStart_withError_setsErrorStatus() { false, new RuntimeException("timed out"))); - plugin.onInvocationEnd(new InvocationEndInfo("req-1", "arn:exec1", true, InvocationStatus.SUCCEEDED, null)); + plugin.onInvocationEnd(new InvocationEndInfo( + "req-1", "arn:exec1", true, Instant.now(), Map.of(), InvocationStatus.SUCCEEDED, null)); var continuationSpan = spanExporter.getFinishedSpanItems().stream() .filter(s -> s.getName().contains("callback")) @@ -1113,14 +1192,14 @@ void operationEnd_withoutMatchingStart_withError_setsErrorStatus() { @Test void contextOperation_doesNotCreateAttemptSpan() { - plugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec1", true, Instant.now())); + plugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec1", true, Instant.now(), Map.of(), Map.of())); // Create operation span first so the CONTEXT user function has a parent plugin.onOperationStart(new OperationInfo( "op-1", "child-ctx", "CONTEXT", "RunInChildContext", null, Instant.now(), null, null, false)); plugin.onUserFunctionStart(new UserFunctionStartInfo( - "op-1", "child-ctx", "CONTEXT", "RunInChildContext", null, Instant.now(), false, null)); + "op-1", "child-ctx", "CONTEXT", "RunInChildContext", null, Instant.now(), false, false, null)); plugin.onUserFunctionEnd(new UserFunctionEndInfo( "op-1", @@ -1131,11 +1210,13 @@ void contextOperation_doesNotCreateAttemptSpan() { Instant.now(), Instant.now(), false, + false, null, false, new SuspendExecutionException())); - plugin.onInvocationEnd(new InvocationEndInfo("req-1", "arn:exec1", true, InvocationStatus.PENDING, null)); + plugin.onInvocationEnd(new InvocationEndInfo( + "req-1", "arn:exec1", true, Instant.now(), Map.of(), InvocationStatus.PENDING, null)); // Should only have the operation span + invocation span — no attempt span var spans = spanExporter.getFinishedSpanItems(); @@ -1148,14 +1229,15 @@ void contextOperation_doesNotCreateAttemptSpan() { @Test void attemptSpan_endedAtInvocationEnd_whenUserFunctionEndNotCalled() { - plugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec1", true, Instant.now())); + plugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec1", true, Instant.now(), Map.of(), Map.of())); // Start attempt but never call onUserFunctionEnd (simulates crash before end hook) plugin.onUserFunctionStart( - new UserFunctionStartInfo("op-1", "running", "STEP", "Step", null, Instant.now(), false, 1)); + new UserFunctionStartInfo("op-1", "running", "STEP", "Step", null, Instant.now(), false, false, 1)); // Invocation ends — attempt span should be cleaned up - plugin.onInvocationEnd(new InvocationEndInfo("req-1", "arn:exec1", true, InvocationStatus.PENDING, null)); + plugin.onInvocationEnd(new InvocationEndInfo( + "req-1", "arn:exec1", true, Instant.now(), Map.of(), InvocationStatus.PENDING, null)); var spans = spanExporter.getFinishedSpanItems(); var attemptSpan = spans.stream() @@ -1169,7 +1251,7 @@ void attemptSpan_endedAtInvocationEnd_whenUserFunctionEndNotCalled() { @Test void childOperation_parentedToParentOperationSpan() { - plugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec1", true, Instant.now())); + plugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec1", true, Instant.now(), Map.of(), Map.of())); // Parent context operation plugin.onOperationStart(new OperationInfo( @@ -1204,7 +1286,8 @@ void childOperation_parentedToParentOperationSpan() { false, null)); - plugin.onInvocationEnd(new InvocationEndInfo("req-1", "arn:exec1", true, InvocationStatus.SUCCEEDED, null)); + plugin.onInvocationEnd(new InvocationEndInfo( + "req-1", "arn:exec1", true, Instant.now(), Map.of(), InvocationStatus.SUCCEEDED, null)); var spans = spanExporter.getFinishedSpanItems(); @@ -1230,18 +1313,19 @@ void multiInvocation_stepWaitStep_producesCorrectSpans() { var arn = "arn:aws:lambda:us-east-1:123:function:test:$LATEST/durable/exec1"; // Invocation 1: step completes, wait starts - plugin.onInvocationStart(new InvocationInfo("req-1", arn, true, Instant.now())); + plugin.onInvocationStart(new InvocationInfo("req-1", arn, true, Instant.now(), Map.of(), Map.of())); plugin.onOperationStart( new OperationInfo("op-1", "step-A", "STEP", "Step", null, Instant.now(), null, null, false)); plugin.onUserFunctionStart( - new UserFunctionStartInfo("op-1", "step-A", "STEP", "Step", null, Instant.now(), false, 1)); + new UserFunctionStartInfo("op-1", "step-A", "STEP", "Step", null, Instant.now(), false, false, 1)); plugin.onUserFunctionEnd(new UserFunctionEndInfo( - "op-1", "step-A", "STEP", "Step", null, Instant.now(), Instant.now(), false, 1, true, null)); + "op-1", "step-A", "STEP", "Step", null, Instant.now(), Instant.now(), false, false, 1, true, null)); plugin.onOperationEnd(new OperationEndInfo( "op-1", "step-A", "STEP", "Step", null, Instant.now(), Instant.now(), "SUCCEEDED", null, false, null)); plugin.onOperationStart( new OperationInfo("op-2", "pause", "WAIT", "Wait", null, Instant.now(), null, null, false)); - plugin.onInvocationEnd(new InvocationEndInfo("req-1", arn, true, InvocationStatus.PENDING, null)); + plugin.onInvocationEnd( + new InvocationEndInfo("req-1", arn, true, Instant.now(), Map.of(), InvocationStatus.PENDING, null)); // Invocation 1 should have: step op + step attempt + wait (PENDING) + invocation = 4 assertEquals(4, spanExporter.getFinishedSpanItems().size()); @@ -1250,18 +1334,19 @@ void multiInvocation_stepWaitStep_producesCorrectSpans() { spanExporter.reset(); // Invocation 2: wait completed between invocations, new step runs - plugin.onInvocationStart(new InvocationInfo("req-2", arn, false, Instant.now())); + plugin.onInvocationStart(new InvocationInfo("req-2", arn, false, Instant.now(), Map.of(), Map.of())); plugin.onOperationEnd(new OperationEndInfo( "op-2", "pause", "WAIT", "Wait", null, Instant.now(), Instant.now(), "SUCCEEDED", null, false, null)); plugin.onOperationStart( new OperationInfo("op-3", "step-B", "STEP", "Step", null, Instant.now(), null, null, false)); plugin.onUserFunctionStart( - new UserFunctionStartInfo("op-3", "step-B", "STEP", "Step", null, Instant.now(), false, 1)); + new UserFunctionStartInfo("op-3", "step-B", "STEP", "Step", null, Instant.now(), false, false, 1)); plugin.onUserFunctionEnd(new UserFunctionEndInfo( - "op-3", "step-B", "STEP", "Step", null, Instant.now(), Instant.now(), false, 1, true, null)); + "op-3", "step-B", "STEP", "Step", null, Instant.now(), Instant.now(), false, false, 1, true, null)); plugin.onOperationEnd(new OperationEndInfo( "op-3", "step-B", "STEP", "Step", null, Instant.now(), Instant.now(), "SUCCEEDED", null, false, null)); - plugin.onInvocationEnd(new InvocationEndInfo("req-2", arn, false, InvocationStatus.SUCCEEDED, null)); + plugin.onInvocationEnd( + new InvocationEndInfo("req-2", arn, false, Instant.now(), Map.of(), InvocationStatus.SUCCEEDED, null)); var inv2Spans = spanExporter.getFinishedSpanItems(); // wait continuation + step-B op + step-B attempt + invocation + Workflow = 5 @@ -1286,11 +1371,11 @@ void crossInvocation_stepRetry_attemptsParentedToRespectiveInvocations() { var arn = "arn:aws:lambda:us-east-1:123:function:test:$LATEST/durable/exec1"; // Invocation 1: step starts, attempt 1 fails, invocation suspended during retry poll - plugin.onInvocationStart(new InvocationInfo("req-1", arn, true, Instant.now())); + plugin.onInvocationStart(new InvocationInfo("req-1", arn, true, Instant.now(), Map.of(), Map.of())); plugin.onOperationStart( new OperationInfo("op-1", "process-payment", "STEP", "Step", null, Instant.now(), null, null, false)); - plugin.onUserFunctionStart( - new UserFunctionStartInfo("op-1", "process-payment", "STEP", "Step", null, Instant.now(), false, 1)); + plugin.onUserFunctionStart(new UserFunctionStartInfo( + "op-1", "process-payment", "STEP", "Step", null, Instant.now(), false, false, 1)); plugin.onUserFunctionEnd(new UserFunctionEndInfo( "op-1", "process-payment", @@ -1300,10 +1385,12 @@ void crossInvocation_stepRetry_attemptsParentedToRespectiveInvocations() { Instant.now(), Instant.now(), false, + false, 1, false, new RuntimeException("payment failed"))); - plugin.onInvocationEnd(new InvocationEndInfo("req-1", arn, true, InvocationStatus.PENDING, null)); + plugin.onInvocationEnd( + new InvocationEndInfo("req-1", arn, true, Instant.now(), Map.of(), InvocationStatus.PENDING, null)); var inv1Spans = spanExporter.getFinishedSpanItems(); // operation span (PENDING) + attempt 1 span + invocation span = 3 @@ -1328,14 +1415,25 @@ void crossInvocation_stepRetry_attemptsParentedToRespectiveInvocations() { spanExporter.reset(); // Invocation 2: step is replayed (continuation), attempt 2 executes and succeeds - plugin.onInvocationStart(new InvocationInfo("req-2", arn, false, Instant.now())); + plugin.onInvocationStart(new InvocationInfo("req-2", arn, false, Instant.now(), Map.of(), Map.of())); // isReplay=true: this operation already exists in the execution state plugin.onOperationStart( new OperationInfo("op-1", "process-payment", "STEP", "Step", null, Instant.now(), null, null, true)); - plugin.onUserFunctionStart( - new UserFunctionStartInfo("op-1", "process-payment", "STEP", "Step", null, Instant.now(), false, 2)); + plugin.onUserFunctionStart(new UserFunctionStartInfo( + "op-1", "process-payment", "STEP", "Step", null, Instant.now(), false, false, 2)); plugin.onUserFunctionEnd(new UserFunctionEndInfo( - "op-1", "process-payment", "STEP", "Step", null, Instant.now(), Instant.now(), false, 2, true, null)); + "op-1", + "process-payment", + "STEP", + "Step", + null, + Instant.now(), + Instant.now(), + false, + false, + 2, + true, + null)); plugin.onOperationEnd(new OperationEndInfo( "op-1", "process-payment", @@ -1348,7 +1446,8 @@ void crossInvocation_stepRetry_attemptsParentedToRespectiveInvocations() { null, false, null)); - plugin.onInvocationEnd(new InvocationEndInfo("req-2", arn, false, InvocationStatus.SUCCEEDED, null)); + plugin.onInvocationEnd( + new InvocationEndInfo("req-2", arn, false, Instant.now(), Map.of(), InvocationStatus.SUCCEEDED, null)); var inv2Spans = spanExporter.getFinishedSpanItems(); // operation span + attempt 2 span + invocation span + Workflow span = 4 @@ -1394,8 +1493,9 @@ void crossInvocation_stepRetry_attemptsParentedToRespectiveInvocations() { @Test void workflowSpan_exportedOnTerminal_internal_deterministicId() { - plugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec-wf", true, Instant.now())); - plugin.onInvocationEnd(new InvocationEndInfo("req-1", "arn:exec-wf", true, InvocationStatus.SUCCEEDED, null)); + plugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec-wf", true, Instant.now(), Map.of(), Map.of())); + plugin.onInvocationEnd(new InvocationEndInfo( + "req-1", "arn:exec-wf", true, Instant.now(), Map.of(), InvocationStatus.SUCCEEDED, null)); var workflow = spanByName("Workflow"); assertEquals(SpanKind.INTERNAL, workflow.getKind(), "Workflow span must be INTERNAL"); @@ -1405,8 +1505,9 @@ void workflowSpan_exportedOnTerminal_internal_deterministicId() { @Test void workflowSpan_notExportedOnNonTerminal() { - plugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec1", true, Instant.now())); - plugin.onInvocationEnd(new InvocationEndInfo("req-1", "arn:exec1", true, InvocationStatus.PENDING, null)); + plugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec1", true, Instant.now(), Map.of(), Map.of())); + plugin.onInvocationEnd(new InvocationEndInfo( + "req-1", "arn:exec1", true, Instant.now(), Map.of(), InvocationStatus.PENDING, null)); assertTrue( spanExporter.getFinishedSpanItems().stream() @@ -1416,16 +1517,17 @@ void workflowSpan_notExportedOnNonTerminal() { @Test void operationAndAttemptSpans_linkToWorkflowSpan() { - plugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec-wf", true, Instant.now())); + plugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec-wf", true, Instant.now(), Map.of(), Map.of())); plugin.onOperationStart( new OperationInfo("op-1", "step-a", "STEP", "Step", null, Instant.now(), null, null, false)); plugin.onUserFunctionStart( - new UserFunctionStartInfo("op-1", "step-a", "STEP", "Step", null, Instant.now(), false, 1)); + new UserFunctionStartInfo("op-1", "step-a", "STEP", "Step", null, Instant.now(), false, false, 1)); plugin.onUserFunctionEnd(new UserFunctionEndInfo( - "op-1", "step-a", "STEP", "Step", null, Instant.now(), Instant.now(), false, 1, true, null)); + "op-1", "step-a", "STEP", "Step", null, Instant.now(), Instant.now(), false, false, 1, true, null)); plugin.onOperationEnd(new OperationEndInfo( "op-1", "step-a", "STEP", "Step", null, Instant.now(), Instant.now(), "SUCCEEDED", 1, false, null)); - plugin.onInvocationEnd(new InvocationEndInfo("req-1", "arn:exec-wf", true, InvocationStatus.SUCCEEDED, null)); + plugin.onInvocationEnd(new InvocationEndInfo( + "req-1", "arn:exec-wf", true, Instant.now(), Map.of(), InvocationStatus.SUCCEEDED, null)); var workflowId = spanByName("Workflow").getSpanId(); var operationSpan = spanByName("step-a"); @@ -1449,12 +1551,13 @@ void operationLinksToWorkflow_withXRayContext() { () -> new ExtractedContext("5759e988bd862e3fe1be46a994272793", "53995c3f42cd8ad8")) .enableMdc(false) .build()); - xrayPlugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec1", true, Instant.now())); + xrayPlugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec1", true, Instant.now(), Map.of(), Map.of())); xrayPlugin.onOperationStart( new OperationInfo("op-1", "step-a", "STEP", "Step", null, Instant.now(), null, null, false)); xrayPlugin.onOperationEnd(new OperationEndInfo( "op-1", "step-a", "STEP", "Step", null, Instant.now(), Instant.now(), "SUCCEEDED", null, false, null)); - xrayPlugin.onInvocationEnd(new InvocationEndInfo("req-1", "arn:exec1", true, InvocationStatus.SUCCEEDED, null)); + xrayPlugin.onInvocationEnd(new InvocationEndInfo( + "req-1", "arn:exec1", true, Instant.now(), Map.of(), InvocationStatus.SUCCEEDED, null)); var workflowId = exporter.getFinishedSpanItems().stream() .filter(s -> s.getName().equals("Workflow")) @@ -1481,9 +1584,10 @@ void workflowSpanName_isConfigurable() { .enableMdc(false) .workflowSpanName("MyWorkflow") .build()); - customPlugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec1", true, Instant.now())); - customPlugin.onInvocationEnd( - new InvocationEndInfo("req-1", "arn:exec1", true, InvocationStatus.SUCCEEDED, null)); + customPlugin.onInvocationStart( + new InvocationInfo("req-1", "arn:exec1", true, Instant.now(), Map.of(), Map.of())); + customPlugin.onInvocationEnd(new InvocationEndInfo( + "req-1", "arn:exec1", true, Instant.now(), Map.of(), InvocationStatus.SUCCEEDED, null)); assertTrue( exporter.getFinishedSpanItems().stream() @@ -1497,9 +1601,15 @@ void workflowSpanName_isConfigurable() { @Test void failedInvocation_setsErrorOnBothWorkflowAndInvocationSpans() { - plugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec1", true, Instant.now())); + plugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec1", true, Instant.now(), Map.of(), Map.of())); plugin.onInvocationEnd(new InvocationEndInfo( - "req-1", "arn:exec1", true, InvocationStatus.FAILED, new RuntimeException("boom"))); + "req-1", + "arn:exec1", + true, + Instant.now(), + Map.of(), + InvocationStatus.FAILED, + new RuntimeException("boom"))); var workflowSpan = spanByName("Workflow"); var invocationSpan = spanByName("Invocation"); diff --git a/otel-plugin/src/test/java/software/amazon/lambda/durable/otel/MdcSpanEnricherTest.java b/otel-plugin/src/test/java/software/amazon/lambda/durable/otel/MdcSpanEnricherTest.java index 2318e4075..708441015 100644 --- a/otel-plugin/src/test/java/software/amazon/lambda/durable/otel/MdcSpanEnricherTest.java +++ b/otel-plugin/src/test/java/software/amazon/lambda/durable/otel/MdcSpanEnricherTest.java @@ -8,6 +8,7 @@ import io.opentelemetry.sdk.trace.SdkTracerProvider; import io.opentelemetry.sdk.trace.export.SimpleSpanProcessor; import java.time.Instant; +import java.util.Map; import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.Test; import org.slf4j.MDC; @@ -63,10 +64,11 @@ void plugin_withMdcEnabled_setsFieldsInMdc() { .enableMdc(true) .build()); - plugin.onInvocationStart(new InvocationInfo("req-1", "arn:exec-mdc-test", true, Instant.now())); + plugin.onInvocationStart( + new InvocationInfo("req-1", "arn:exec-mdc-test", true, Instant.now(), Map.of(), Map.of())); plugin.onUserFunctionStart( - new UserFunctionStartInfo("op-1", "step", "STEP", "Step", null, Instant.now(), false, 1)); + new UserFunctionStartInfo("op-1", "step", "STEP", "Step", null, Instant.now(), false, false, 1)); // MDC should have trace fields after onUserFunctionStart assertNotNull(MDC.get(MdcSpanEnricher.MDC_TRACE_ID)); @@ -74,15 +76,15 @@ void plugin_withMdcEnabled_setsFieldsInMdc() { assertNotNull(MDC.get(MdcSpanEnricher.MDC_TRACE_SAMPLED)); plugin.onUserFunctionEnd(new UserFunctionEndInfo( - "op-1", "step", "STEP", "Step", null, Instant.now(), Instant.now(), false, 1, true, null)); + "op-1", "step", "STEP", "Step", null, Instant.now(), Instant.now(), false, false, 1, true, null)); // After onUserFunctionEnd: span_id is cleared, but trace_id remains for handler-level logs between steps assertNotNull(MDC.get(MdcSpanEnricher.MDC_TRACE_ID), "trace_id should persist between steps"); assertNull(MDC.get(MdcSpanEnricher.MDC_SPAN_ID), "span_id should be cleared after step"); assertNotNull(MDC.get(MdcSpanEnricher.MDC_TRACE_SAMPLED), "trace_flags should persist between steps"); - plugin.onInvocationEnd( - new InvocationEndInfo("req-1", "arn:exec-mdc-test", true, InvocationStatus.SUCCEEDED, null)); + plugin.onInvocationEnd(new InvocationEndInfo( + "req-1", "arn:exec-mdc-test", true, Instant.now(), Map.of(), InvocationStatus.SUCCEEDED, null)); // After onInvocationEnd: all MDC fields are cleared assertNull(MDC.get(MdcSpanEnricher.MDC_TRACE_ID)); diff --git a/sdk/src/main/java/software/amazon/lambda/durable/execution/DurableExecutor.java b/sdk/src/main/java/software/amazon/lambda/durable/execution/DurableExecutor.java index 6ab6bdbc3..617ed4639 100644 --- a/sdk/src/main/java/software/amazon/lambda/durable/execution/DurableExecutor.java +++ b/sdk/src/main/java/software/amazon/lambda/durable/execution/DurableExecutor.java @@ -28,6 +28,7 @@ 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.PluginInfoConverter; import software.amazon.lambda.durable.plugin.PluginRunner; import software.amazon.lambda.durable.serde.SerDes; import software.amazon.lambda.durable.util.ExceptionHelper; @@ -67,11 +68,19 @@ public static DurableExecutionOutput execute( // onInvocationStart runs on the user thread so plugins can // inject ThreadLocal objects, update MDC, etc. // executionStartTime comes from the initial EXECUTION operation in the first backend event. + // The operation maps are snapshots of the state delivered for this invocation; on a replay + // invocation updatedOperations names the operations the backend completed while suspended. pluginRunner.onInvocationStart(new InvocationInfo( requestId, executionArn, isFirstInvocation, - executionManager.getExecutionOperation().startTimestamp())); + executionManager.getExecutionOperation().startTimestamp(), + PluginInfoConverter.toOperationItemMap( + executionManager.getOperationsSnapshot(), + executionManager.getInitialOperationIds()), + PluginInfoConverter.toOperationItemMap( + executionManager.getUpdatedOperationsSnapshot(), + executionManager.getInitialOperationIds()))); var userInput = extractUserInput( executionManager.getExecutionOperation(), config.getSerDes(), inputType); @@ -99,6 +108,7 @@ public static DurableExecutionOutput execute( if (cause instanceof SuspendExecutionException) { fireOnInvocationEnd( pluginRunner, + executionManager, requestId, executionArn, isFirstInvocation, @@ -115,6 +125,7 @@ public static DurableExecutionOutput execute( && unrecoverableDurableExecutionException.isRetryable()) { fireOnInvocationEnd( pluginRunner, + executionManager, requestId, executionArn, isFirstInvocation, @@ -127,6 +138,7 @@ public static DurableExecutionOutput execute( logger.debug("Execution failed: {}", cause.getMessage()); fireOnInvocationEnd( pluginRunner, + executionManager, requestId, executionArn, isFirstInvocation, @@ -141,6 +153,7 @@ public static DurableExecutionOutput execute( DurableExecutionOutput.success(handleLargePayload(executionManager, outputPayload)); fireOnInvocationEnd( pluginRunner, + executionManager, requestId, executionArn, isFirstInvocation, @@ -159,12 +172,24 @@ public static DurableExecutionOutput execute( private static void fireOnInvocationEnd( PluginRunner pluginRunner, + ExecutionManager executionManager, String requestId, String executionArn, boolean isFirstInvocation, InvocationStatus status, Throwable error) { - pluginRunner.onInvocationEnd(new InvocationEndInfo(requestId, executionArn, isFirstInvocation, status, error)); + // The end info repeats the start info's identity surface (execution start time, operation snapshot) so an + // invocation-end hook never has to correlate back to the start hook. The snapshot is taken at end time, so + // unlike the start info it also contains operations created during this invocation. + pluginRunner.onInvocationEnd(new InvocationEndInfo( + requestId, + executionArn, + isFirstInvocation, + executionManager.getExecutionOperation().startTimestamp(), + PluginInfoConverter.toOperationItemMap( + executionManager.getOperationsSnapshot(), executionManager.getInitialOperationIds()), + status, + error)); } private static String handleLargePayload(ExecutionManager executionManager, String outputPayload) { diff --git a/sdk/src/main/java/software/amazon/lambda/durable/execution/ExecutionManager.java b/sdk/src/main/java/software/amazon/lambda/durable/execution/ExecutionManager.java index 1c45cb0d6..b52216d99 100644 --- a/sdk/src/main/java/software/amazon/lambda/durable/execution/ExecutionManager.java +++ b/sdk/src/main/java/software/amazon/lambda/durable/execution/ExecutionManager.java @@ -5,6 +5,7 @@ import com.amazonaws.services.lambda.runtime.Context; import java.time.Instant; import java.util.ArrayList; +import java.util.Collection; import java.util.Collections; import java.util.HashSet; import java.util.List; @@ -63,6 +64,7 @@ public class ExecutionManager implements SafeCloseable { private final AtomicReference executionMode; private final DurableConfig durableConfig; private final Set updatedOperationIdsSinceLastInvocation; + private final Set initialOperationIds; // ===== Thread Coordination ===== private final Map registeredOperations = new ConcurrentHashMap<>(); @@ -89,6 +91,11 @@ public ExecutionManager(DurableExecutionInput input, DurableConfig config, Conte this.operationStorage = checkpointManager.fetchAllPages(input.initialExecutionState()).stream() .collect(Collectors.toConcurrentMap(Operation::id, op -> op)); + // The ids delivered in this invocation's initial state. Everything else in operationStorage is created during + // this invocation, so this set is what distinguishes replayed operations from freshly-started ones for the + // plugin hooks' isReplay indicators. + this.initialOperationIds = Set.copyOf(operationStorage.keySet()); + // Start in REPLAY mode if we have more than just the initial EXECUTION operation this.executionMode = new AtomicReference<>(operationStorage.size() > 1 ? ExecutionMode.REPLAY : ExecutionMode.EXECUTION); @@ -132,6 +139,47 @@ public boolean isOperationUpdatedSinceLastInvocation(String operationId) { return updatedOperationIdsSinceLastInvocation.contains(operationId); } + /** + * Returns {@code true} if the given operation was present in the checkpointed state delivered at the start of this + * invocation, i.e. it predates this invocation and is being replayed rather than started fresh. Unlike + * {@link #getOperationAndUpdateReplayState(String)} this does not mutate the execution's replay mode, so it is safe + * to call from plugin-hook firing sites. + * + * @param operationId the operation ID to check + * @return true if the operation was delivered in this invocation's initial state + */ + public boolean wasObservedAtInvocationStart(String operationId) { + return initialOperationIds.contains(operationId); + } + + /** Returns the ids of the operations delivered in this invocation's initial state. */ + public Set getInitialOperationIds() { + return initialOperationIds; + } + + /** + * Returns an immutable snapshot of the operations currently tracked for this execution, including the initial + * EXECUTION operation. Non-mutating; intended for the invocation-level plugin hooks. + * + * @return a snapshot of the tracked operations + */ + public Collection getOperationsSnapshot() { + return List.copyOf(operationStorage.values()); + } + + /** + * Returns the subset of {@link #getOperationsSnapshot()} whose ids the backend reported as updated since the last + * successful invocation. Empty on the first invocation. Ids without a corresponding tracked operation are skipped. + * + * @return a snapshot of the externally-updated operations + */ + public Collection getUpdatedOperationsSnapshot() { + return updatedOperationIdsSinceLastInvocation.stream() + .map(operationStorage::get) + .filter(Objects::nonNull) + .toList(); + } + /** Registers an operation so it can receive checkpoint completion notifications. */ public void registerOperation(BaseDurableOperation operation) { registeredOperations.put(operation.getOperationId(), operation); @@ -162,7 +210,11 @@ private void onCheckpointComplete(List newOperations) { durableConfig .getPluginRunner() .onOperationChange(PluginInfoConverter.toOperationChangeInfo( - requestId, durableExecutionArn, updatedOperations, operationStorage.values())); + requestId, + durableExecutionArn, + updatedOperations, + operationStorage.values(), + initialOperationIds)); } } diff --git a/sdk/src/main/java/software/amazon/lambda/durable/operation/BaseDurableOperation.java b/sdk/src/main/java/software/amazon/lambda/durable/operation/BaseDurableOperation.java index f0391567a..9a6e92771 100644 --- a/sdk/src/main/java/software/amazon/lambda/durable/operation/BaseDurableOperation.java +++ b/sdk/src/main/java/software/amazon/lambda/durable/operation/BaseDurableOperation.java @@ -350,7 +350,11 @@ protected void runUserHandler(Runnable runnable, ThreadType threadType) { protected T runUserFunction(Integer attempt, Supplier userFunction) { var pluginRunner = getPluginRunner(); var startInfo = PluginInfoConverter.toUserFunctionStartInfo( - operationIdentifier, durableContext.getParentId(), durableContext.isReplaying(), attempt); + operationIdentifier, + durableContext.getParentId(), + executionManager.wasObservedAtInvocationStart(getOperationId()), + durableContext.isReplaying(), + attempt); pluginRunner.onUserFunctionStart(startInfo); try { T result = userFunction.get(); diff --git a/sdk/src/main/java/software/amazon/lambda/durable/plugin/InvocationEndInfo.java b/sdk/src/main/java/software/amazon/lambda/durable/plugin/InvocationEndInfo.java index ba85b9e70..3e48ec061 100644 --- a/sdk/src/main/java/software/amazon/lambda/durable/plugin/InvocationEndInfo.java +++ b/sdk/src/main/java/software/amazon/lambda/durable/plugin/InvocationEndInfo.java @@ -2,12 +2,23 @@ // SPDX-License-Identifier: Apache-2.0 package software.amazon.lambda.durable.plugin; +import java.time.Instant; +import java.util.Map; + /** * Information provided at the end of a Lambda invocation. * + *

Carries the same invocation-identity surface as {@link InvocationInfo} so an invocation-end hook never has to + * correlate back to the start hook to learn the execution start time or the operation state. + * * @param requestId the Lambda request ID for this invocation * @param durableExecutionArn the durable execution ARN * @param isFirstInvocation true if this is the first invocation of the execution + * @param executionStartTime the start timestamp of the durable execution, taken from the initial EXECUTION operation in + * the first event delivered by the backend. Stable across all invocations of the same execution. + * @param operations a snapshot of the checkpointed operations known when the invocation ended, keyed by operation ID. + * Unlike {@link InvocationInfo#operations()} this includes operations created during this invocation. Includes the + * initial EXECUTION operation. Empty-but-never-null. * @param invocationStatus the invocation outcome (SUCCEEDED, FAILED, or PENDING) * @param executionError non-null if the execution failed * @deprecated This is a preview API that is experimental and may be changed or removed in future releases. @@ -17,5 +28,7 @@ public record InvocationEndInfo( String requestId, String durableExecutionArn, boolean isFirstInvocation, + Instant executionStartTime, + Map operations, InvocationStatus invocationStatus, Throwable executionError) {} diff --git a/sdk/src/main/java/software/amazon/lambda/durable/plugin/InvocationInfo.java b/sdk/src/main/java/software/amazon/lambda/durable/plugin/InvocationInfo.java index 8e2406829..9fbb38c22 100644 --- a/sdk/src/main/java/software/amazon/lambda/durable/plugin/InvocationInfo.java +++ b/sdk/src/main/java/software/amazon/lambda/durable/plugin/InvocationInfo.java @@ -3,6 +3,7 @@ package software.amazon.lambda.durable.plugin; import java.time.Instant; +import java.util.Map; /** * Invocation-level information available to plugin hooks. @@ -12,8 +13,19 @@ * @param isFirstInvocation true if this is the first invocation of the execution (not a replay invocation) * @param executionStartTime the start timestamp of the durable execution, taken from the initial EXECUTION operation in * the first event delivered by the backend. Stable across all invocations of the same execution. + * @param operations a snapshot of the checkpointed operations delivered at the start of this invocation, keyed by + * operation ID. Includes the initial EXECUTION operation. Empty-but-never-null. + * @param updatedOperations the subset of {@code operations} that changed externally between the previous invocation and + * this one (a wait timer expired, a callback was received, a chained invoke completed), keyed by operation ID. + * Sourced from the {@code UpdatedOperationIds} field of the durable invocation input, so it is empty on the first + * invocation. Empty-but-never-null. * @deprecated This is a preview API that is experimental and may be changed or removed in future releases. */ @Deprecated public record InvocationInfo( - String requestId, String durableExecutionArn, boolean isFirstInvocation, Instant executionStartTime) {} + String requestId, + String durableExecutionArn, + boolean isFirstInvocation, + Instant executionStartTime, + Map operations, + Map updatedOperations) {} diff --git a/sdk/src/main/java/software/amazon/lambda/durable/plugin/OperationChangeItemInfo.java b/sdk/src/main/java/software/amazon/lambda/durable/plugin/OperationChangeItemInfo.java index 6846280fe..c33de47a3 100644 --- a/sdk/src/main/java/software/amazon/lambda/durable/plugin/OperationChangeItemInfo.java +++ b/sdk/src/main/java/software/amazon/lambda/durable/plugin/OperationChangeItemInfo.java @@ -6,7 +6,11 @@ import software.amazon.awssdk.services.lambda.model.OperationStatus; /** - * Operation-level information for a single operation within an {@link OperationChangeInfo}. + * Operation-level information for a single operation within an {@link OperationChangeInfo}, and the snapshot record + * used for the operation maps carried on {@link InvocationInfo} / {@link InvocationEndInfo}. + * + *

Carries the full operation field surface, mirroring {@link OperationEndInfo}, so a plugin observing an operation + * through a change delta or an invocation-level map sees the same fields it would see through the per-operation hooks. * * @param id operation ID * @param name human-readable operation name (may be null) @@ -15,8 +19,11 @@ * @param parentId parent operation ID (null for root-level operations) * @param startTimestamp when the operation started * @param endTimestamp when the operation ended - * @param error non-null if the operation failed * @param status operation status + * @param attempt the attempt number for retriable operations (STEP, WAIT_FOR_CONDITION) — null for others + * @param isReplay true if this operation was already present in the checkpointed state delivered at the start of the + * current invocation (i.e. it predates this invocation) rather than being created during it + * @param error non-null if the operation failed * @deprecated This is a preview API that is experimental and may be changed or removed in future releases. */ @Deprecated @@ -28,5 +35,7 @@ public record OperationChangeItemInfo( String parentId, Instant startTimestamp, Instant endTimestamp, - Throwable error, - OperationStatus status) {} + OperationStatus status, + Integer attempt, + boolean isReplay, + Throwable error) {} diff --git a/sdk/src/main/java/software/amazon/lambda/durable/plugin/PluginInfoConverter.java b/sdk/src/main/java/software/amazon/lambda/durable/plugin/PluginInfoConverter.java index 58911f7be..57433415d 100644 --- a/sdk/src/main/java/software/amazon/lambda/durable/plugin/PluginInfoConverter.java +++ b/sdk/src/main/java/software/amazon/lambda/durable/plugin/PluginInfoConverter.java @@ -4,6 +4,8 @@ import java.time.Instant; import java.util.Collection; +import java.util.Map; +import java.util.Set; import java.util.stream.Collectors; import software.amazon.awssdk.services.lambda.model.Operation; import software.amazon.lambda.durable.model.OperationIdentifier; @@ -77,12 +79,17 @@ public static OperationEndInfo toOperationEndInfo( * * @param identifier the operation identifier containing id, name, type, and subType * @param parentId the parent operation ID (may be null) - * @param isReplay true if the user function is called during replay (context operations) + * @param isReplay true if this operation was already present in the checkpointed state when it started + * @param isReplayingChildren true if the child operations of this context body are replaying from checkpoints * @param attempt the 1-based attempt number (null for context operations) * @return a UserFunctionStartInfo record */ public static UserFunctionStartInfo toUserFunctionStartInfo( - OperationIdentifier identifier, String parentId, boolean isReplayingChildren, Integer attempt) { + OperationIdentifier identifier, + String parentId, + boolean isReplay, + boolean isReplayingChildren, + Integer attempt) { return new UserFunctionStartInfo( identifier.operationId(), identifier.name(), @@ -90,6 +97,7 @@ public static UserFunctionStartInfo toUserFunctionStartInfo( identifier.subType() != null ? identifier.subType().getValue() : null, parentId, Instant.now(), + isReplay, isReplayingChildren, attempt); } @@ -112,6 +120,7 @@ public static UserFunctionEndInfo toUserFunctionEndInfo( startInfo.parentId(), startInfo.startTimestamp(), Instant.now(), + startInfo.isReplay(), startInfo.isReplayingChildren(), startInfo.attempt(), succeeded, @@ -126,25 +135,41 @@ public static UserFunctionEndInfo toUserFunctionEndInfo( * @param durableExecutionArn the durable execution ARN * @param updatedOperations the durable operations whose status changed in this checkpoint response * @param allOperations all durable operations tracked for the execution after this response + * @param replayedOperationIds ids of the operations delivered in this invocation's initial state, used to populate + * each item's {@code isReplay} indicator * @return an OperationChangeInfo record */ public static OperationChangeInfo toOperationChangeInfo( String requestId, String durableExecutionArn, Collection updatedOperations, - Collection allOperations) { + Collection allOperations, + Set replayedOperationIds) { return new OperationChangeInfo( requestId, durableExecutionArn, - updatedOperations.stream() - .collect(Collectors.toUnmodifiableMap( - Operation::id, PluginInfoConverter::toOperationChangeItemInfo)), - allOperations.stream() - .collect(Collectors.toUnmodifiableMap( - Operation::id, PluginInfoConverter::toOperationChangeItemInfo))); + toOperationItemMap(updatedOperations, replayedOperationIds), + toOperationItemMap(allOperations, replayedOperationIds)); } - private static OperationChangeItemInfo toOperationChangeItemInfo(Operation operation) { + /** + * Converts durable operations to an unmodifiable map of {@link OperationChangeItemInfo}, keyed by operation ID. + * + * @param operations the durable operations to convert + * @param replayedOperationIds ids of the operations delivered in this invocation's initial state, used to populate + * each item's {@code isReplay} indicator + * @return an unmodifiable map of operation ID to item info + */ + public static Map toOperationItemMap( + Collection operations, Set replayedOperationIds) { + return operations.stream() + .collect(Collectors.toUnmodifiableMap( + Operation::id, + operation -> + toOperationChangeItemInfo(operation, replayedOperationIds.contains(operation.id())))); + } + + private static OperationChangeItemInfo toOperationChangeItemInfo(Operation operation, boolean isReplay) { return new OperationChangeItemInfo( operation.id(), operation.name(), @@ -153,7 +178,9 @@ private static OperationChangeItemInfo toOperationChangeItemInfo(Operation opera operation.parentId(), operation.startTimestamp(), operation.endTimestamp(), - BaseDurableOperation.extractErrorFromOperation(operation), - operation.status()); + operation.status(), + operation.stepDetails() != null ? operation.stepDetails().attempt() : null, + isReplay, + BaseDurableOperation.extractErrorFromOperation(operation)); } } diff --git a/sdk/src/main/java/software/amazon/lambda/durable/plugin/UserFunctionEndInfo.java b/sdk/src/main/java/software/amazon/lambda/durable/plugin/UserFunctionEndInfo.java index 86136b3bd..63b8b8ca1 100644 --- a/sdk/src/main/java/software/amazon/lambda/durable/plugin/UserFunctionEndInfo.java +++ b/sdk/src/main/java/software/amazon/lambda/durable/plugin/UserFunctionEndInfo.java @@ -16,6 +16,9 @@ * @param parentId parent operation ID (null for root-level operations) * @param startTimestamp when the user function started * @param endTimestamp when the user function ended + * @param isReplay true if THIS operation was already present in the execution's checkpointed state when it started + * (i.e. observed via replay rather than created fresh in this invocation). Distinct from + * {@code isReplayingChildren}, which is about the child operations of a context body. * @param isReplayingChildren true if child operations within this context are being replayed from checkpoints * @param attempt 1-based attempt number for steps/waitForCondition, null for context operations * @param succeeded true if the user function completed without error @@ -31,6 +34,7 @@ public record UserFunctionEndInfo( String parentId, Instant startTimestamp, Instant endTimestamp, + boolean isReplay, boolean isReplayingChildren, Integer attempt, boolean succeeded, diff --git a/sdk/src/main/java/software/amazon/lambda/durable/plugin/UserFunctionStartInfo.java b/sdk/src/main/java/software/amazon/lambda/durable/plugin/UserFunctionStartInfo.java index 842f82536..4285a5026 100644 --- a/sdk/src/main/java/software/amazon/lambda/durable/plugin/UserFunctionStartInfo.java +++ b/sdk/src/main/java/software/amazon/lambda/durable/plugin/UserFunctionStartInfo.java @@ -16,6 +16,9 @@ * @param subType operation sub-type (Map, Parallel, WaitForCondition, etc.) — may be null * @param parentId parent operation ID (null for root-level operations) * @param startTimestamp when the user function started + * @param isReplay true if THIS operation was already present in the execution's checkpointed state when it started + * (i.e. observed via replay rather than created fresh in this invocation). Distinct from + * {@code isReplayingChildren}, which is about the child operations of a context body. * @param isReplayingChildren true if child operations within this context are being replayed from checkpoints * @param attempt 1-based attempt number for steps/waitForCondition, null for context operations * @deprecated This is a preview API that is experimental and may be changed or removed in future releases. @@ -28,5 +31,6 @@ public record UserFunctionStartInfo( String subType, String parentId, Instant startTimestamp, + boolean isReplay, boolean isReplayingChildren, Integer attempt) {} diff --git a/sdk/src/test/java/software/amazon/lambda/durable/DurableConfigTest.java b/sdk/src/test/java/software/amazon/lambda/durable/DurableConfigTest.java index c0bd54147..266e43a6e 100644 --- a/sdk/src/test/java/software/amazon/lambda/durable/DurableConfigTest.java +++ b/sdk/src/test/java/software/amazon/lambda/durable/DurableConfigTest.java @@ -547,7 +547,7 @@ void testBuilder_WithMultiplePlugins_AllRegistered() { config.getPluginRunner() .onInvocationStart(new software.amazon.lambda.durable.plugin.InvocationInfo( - "req-1", "arn:test", true, java.time.Instant.now())); + "req-1", "arn:test", true, java.time.Instant.now(), java.util.Map.of(), java.util.Map.of())); assertEquals(List.of("p1:onInvocationStart", "p2:onInvocationStart"), calls); } @@ -567,7 +567,7 @@ void testBuilder_WithPlugins_CalledMultipleTimes_Replaces() { config.getPluginRunner() .onInvocationStart(new software.amazon.lambda.durable.plugin.InvocationInfo( - "req-1", "arn:test", true, java.time.Instant.now())); + "req-1", "arn:test", true, java.time.Instant.now(), java.util.Map.of(), java.util.Map.of())); assertEquals(List.of("p2:onInvocationStart", "p3:onInvocationStart"), calls); } diff --git a/sdk/src/test/java/software/amazon/lambda/durable/plugin/PluginInfoConverterTest.java b/sdk/src/test/java/software/amazon/lambda/durable/plugin/PluginInfoConverterTest.java index dbba0870f..82724788f 100644 --- a/sdk/src/test/java/software/amazon/lambda/durable/plugin/PluginInfoConverterTest.java +++ b/sdk/src/test/java/software/amazon/lambda/durable/plugin/PluginInfoConverterTest.java @@ -100,7 +100,7 @@ void toOperationEndInfo_nullError_forSuccess() { @Test void toUserFunctionStartInfo_stepAttempt() { - var info = PluginInfoConverter.toUserFunctionStartInfo(STEP_IDENTIFIER, PARENT_ID, false, 3); + var info = PluginInfoConverter.toUserFunctionStartInfo(STEP_IDENTIFIER, PARENT_ID, true, false, 3); assertEquals(OPERATION_ID, info.id()); assertEquals(OPERATION_NAME, info.name()); @@ -108,16 +108,18 @@ void toUserFunctionStartInfo_stepAttempt() { assertEquals("Step", info.subType()); assertEquals(PARENT_ID, info.parentId()); assertNotNull(info.startTimestamp()); + assertTrue(info.isReplay()); assertFalse(info.isReplayingChildren()); assertEquals(3, info.attempt()); } @Test void toUserFunctionStartInfo_contextOperation() { - var info = PluginInfoConverter.toUserFunctionStartInfo(MAP_IDENTIFIER, PARENT_ID, true, null); + var info = PluginInfoConverter.toUserFunctionStartInfo(MAP_IDENTIFIER, PARENT_ID, false, true, null); assertEquals("CONTEXT", info.type()); assertEquals("Map", info.subType()); + assertFalse(info.isReplay()); assertTrue(info.isReplayingChildren()); assertNull(info.attempt()); } @@ -126,7 +128,7 @@ void toUserFunctionStartInfo_contextOperation() { @Test void toUserFunctionEndInfo_succeeded() { - var startInfo = PluginInfoConverter.toUserFunctionStartInfo(STEP_IDENTIFIER, PARENT_ID, false, 1); + var startInfo = PluginInfoConverter.toUserFunctionStartInfo(STEP_IDENTIFIER, PARENT_ID, true, false, 1); var endInfo = PluginInfoConverter.toUserFunctionEndInfo(startInfo, true, null); @@ -134,6 +136,7 @@ void toUserFunctionEndInfo_succeeded() { assertEquals(OPERATION_NAME, endInfo.name()); assertEquals(startInfo.startTimestamp(), endInfo.startTimestamp()); assertNotNull(endInfo.endTimestamp()); + assertTrue(endInfo.isReplay()); assertFalse(endInfo.isReplayingChildren()); assertEquals(1, endInfo.attempt()); assertTrue(endInfo.succeeded()); @@ -143,7 +146,7 @@ void toUserFunctionEndInfo_succeeded() { @Test void toUserFunctionEndInfo_failed() { var error = new RuntimeException("step failed"); - var startInfo = PluginInfoConverter.toUserFunctionStartInfo(STEP_IDENTIFIER, null, false, 2); + var startInfo = PluginInfoConverter.toUserFunctionStartInfo(STEP_IDENTIFIER, null, false, false, 2); var endInfo = PluginInfoConverter.toUserFunctionEndInfo(startInfo, false, error); diff --git a/sdk/src/test/java/software/amazon/lambda/durable/plugin/PluginRunnerTest.java b/sdk/src/test/java/software/amazon/lambda/durable/plugin/PluginRunnerTest.java index d34b5a8ca..1f287e737 100644 --- a/sdk/src/test/java/software/amazon/lambda/durable/plugin/PluginRunnerTest.java +++ b/sdk/src/test/java/software/amazon/lambda/durable/plugin/PluginRunnerTest.java @@ -137,12 +137,23 @@ void pluginRunner_isImmutable() { private static InvocationInfo invocationInfo() { return new InvocationInfo( - "req-123", "arn:aws:lambda:us-east-1:123456789012:function:test", false, Instant.now()); + "req-123", + "arn:aws:lambda:us-east-1:123456789012:function:test", + false, + Instant.now(), + Map.of(), + Map.of()); } private static InvocationEndInfo invocationEndInfo() { return new InvocationEndInfo( - "req-123", "arn:aws:lambda:us-east-1:123456789012:function:test", false, null, null); + "req-123", + "arn:aws:lambda:us-east-1:123456789012:function:test", + false, + Instant.now(), + Map.of(), + null, + null); } private static OperationInfo operationInfo() { @@ -160,12 +171,12 @@ private static OperationChangeInfo operationChangeInfo() { } private static UserFunctionStartInfo attemptInfo() { - return new UserFunctionStartInfo("op-1", "test-step", "STEP", null, null, Instant.now(), false, 1); + return new UserFunctionStartInfo("op-1", "test-step", "STEP", null, null, Instant.now(), false, false, 1); } private static UserFunctionEndInfo attemptEndInfo() { return new UserFunctionEndInfo( - "op-1", "test-step", "STEP", null, null, Instant.now(), Instant.now(), false, 1, true, null); + "op-1", "test-step", "STEP", null, null, Instant.now(), Instant.now(), false, false, 1, true, null); } // ─── Test plugin implementations ─────────────────────────────────────