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 ─────────────────────────────────────