From fdbfe492ab6da2298ce7d5fad10435a7a1ba1887 Mon Sep 17 00:00:00 2001 From: Ayushi Ahjolia Date: Wed, 5 Aug 2026 16:35:06 -0700 Subject: [PATCH] feat(plugin): Add result accessor to OperationEndInfo --- .../durable/otel/ExecutionOtelPluginTest.java | 137 +++++++++- .../otel/InvocationOtelPluginTest.java | 234 ++++++++++++++++-- .../lambda/durable/PluginIntegrationTest.java | 19 ++ .../durable/plugin/OperationEndInfo.java | 4 +- .../durable/plugin/PluginInfoConverter.java | 28 ++- .../plugin/PluginInfoConverterTest.java | 31 +++ .../durable/plugin/PluginRunnerTest.java | 13 +- 7 files changed, 433 insertions(+), 33 deletions(-) diff --git a/otel-plugin/src/test/java/software/amazon/lambda/durable/otel/ExecutionOtelPluginTest.java b/otel-plugin/src/test/java/software/amazon/lambda/durable/otel/ExecutionOtelPluginTest.java index 6936065c2..4a2af8bab 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 @@ -75,7 +75,18 @@ void defaultConstructor_usesGlobalSdkTracerProviderDirectly() { 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)); + "op-1", + "step", + "STEP", + "Step", + null, + Instant.now(), + Instant.now(), + "SUCCEEDED", + null, + false, + null, + null)); defaultPlugin.onInvocationEnd( new InvocationEndInfo("req-1", "arn:exec1", true, InvocationStatus.SUCCEEDED, null)); @@ -230,7 +241,18 @@ void operationSpan_carriesAttemptNumberAtEnd() { 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)); + "op-1", + "flaky", + "STEP", + "Step", + null, + Instant.now(), + Instant.now(), + "SUCCEEDED", + 3, + false, + null, + null)); plugin.onInvocationEnd(new InvocationEndInfo("req-1", ARN, true, InvocationStatus.SUCCEEDED, null)); var operationSpan = spanByName(spanExporter.getFinishedSpanItems(), "flaky"); @@ -247,7 +269,18 @@ void continuationOperationSpan_carriesAttemptNumber() { plugin.onInvocationStart(new InvocationInfo("req-2", ARN, false, Instant.now())); // 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)); + "op-1", + "flaky", + "STEP", + "Step", + null, + Instant.now(), + Instant.now(), + "SUCCEEDED", + 2, + false, + null, + null)); plugin.onInvocationEnd(new InvocationEndInfo("req-2", ARN, false, InvocationStatus.SUCCEEDED, null)); var operationSpan = spanByName(spanExporter.getFinishedSpanItems(), "flaky"); @@ -266,7 +299,7 @@ void operationSpan_startsAtOperationStartTimestamp() { plugin.onInvocationStart(new InvocationInfo("req-1", ARN, true, Instant.now())); 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)); + "op-1", "step-a", "STEP", "Step", null, opStart, opEnd, "SUCCEEDED", null, false, null, null)); plugin.onInvocationEnd(new InvocationEndInfo("req-1", ARN, true, InvocationStatus.SUCCEEDED, null)); var operationSpan = spanByName(spanExporter.getFinishedSpanItems(), "step-a"); @@ -282,7 +315,18 @@ void operationSpan_parentedToWorkflow_linkedToInvocation() { 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)); + "op-1", + "step-a", + "STEP", + "Step", + null, + Instant.now(), + Instant.now(), + "SUCCEEDED", + null, + false, + null, + null)); plugin.onInvocationEnd(new InvocationEndInfo("req-1", ARN, true, InvocationStatus.SUCCEEDED, null)); var spans = spanExporter.getFinishedSpanItems(); @@ -310,7 +354,18 @@ void attemptSpan_childOfOperation_linkedToInvocation() { plugin.onUserFunctionEnd(new UserFunctionEndInfo( "op-1", "compute", "STEP", "Step", null, Instant.now(), Instant.now(), false, 1, true, null)); plugin.onOperationEnd(new OperationEndInfo( - "op-1", "compute", "STEP", "Step", null, Instant.now(), Instant.now(), "SUCCEEDED", null, false, null)); + "op-1", + "compute", + "STEP", + "Step", + null, + Instant.now(), + Instant.now(), + "SUCCEEDED", + null, + false, + null, + null)); plugin.onInvocationEnd(new InvocationEndInfo("req-1", ARN, true, InvocationStatus.SUCCEEDED, null)); var spans = spanExporter.getFinishedSpanItems(); @@ -350,6 +405,7 @@ void childOperation_parentedToParentOperationSpan() { "SUCCEEDED", null, false, + null, null)); plugin.onOperationEnd(new OperationEndInfo( "op-parent", @@ -362,6 +418,7 @@ void childOperation_parentedToParentOperationSpan() { "SUCCEEDED", null, false, + null, null)); plugin.onInvocationEnd(new InvocationEndInfo("req-1", ARN, true, InvocationStatus.SUCCEEDED, null)); @@ -424,7 +481,18 @@ void operationSuccess_setsOkOnOperationSpan() { 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)); + "op-1", + "step-ok", + "STEP", + "Step", + null, + Instant.now(), + Instant.now(), + "SUCCEEDED", + null, + false, + null, + null)); plugin.onInvocationEnd(new InvocationEndInfo("req-1", ARN, true, InvocationStatus.SUCCEEDED, null)); var operationSpan = spanByName(spanExporter.getFinishedSpanItems(), "step-ok"); @@ -450,6 +518,7 @@ void operationEnd_withNonSuccessStatusAndNoError_leavesOperationSpanUnset() { "CANCELLED", null, false, + null, null)); plugin.onInvocationEnd(new InvocationEndInfo("req-1", ARN, true, InvocationStatus.SUCCEEDED, null)); @@ -473,6 +542,7 @@ void operationEnd_withoutStart_nonSuccessStatusAndNoError_leavesContinuationSpan "TIMED_OUT", null, false, + null, null)); plugin.onInvocationEnd(new InvocationEndInfo("req-2", ARN, false, InvocationStatus.SUCCEEDED, null)); @@ -488,7 +558,18 @@ void operationEnd_withNullStatusAndNoError_setsOkOnOperationSpan() { 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)); + "op-ctx", + "my-ctx", + "CONTEXT", + null, + null, + Instant.now(), + Instant.now(), + null, + null, + false, + null, + null)); plugin.onInvocationEnd(new InvocationEndInfo("req-1", ARN, true, InvocationStatus.SUCCEEDED, null)); var operationSpan = spanByName(spanExporter.getFinishedSpanItems(), "my-ctx"); @@ -528,7 +609,18 @@ void operationOpenedThenCompletedNextInvocation_exportedOnceOnOperationEnd() { // Invocation 2: the operation completes → materialized once via onOperationEnd, linked to this invocation. plugin.onInvocationStart(new InvocationInfo("req-2", ARN, false, Instant.now())); plugin.onOperationEnd(new OperationEndInfo( - "op-1", "my-wait", "WAIT", "Wait", null, Instant.now(), Instant.now(), "SUCCEEDED", null, false, null)); + "op-1", + "my-wait", + "WAIT", + "Wait", + null, + Instant.now(), + Instant.now(), + "SUCCEEDED", + null, + false, + null, + null)); plugin.onInvocationEnd(new InvocationEndInfo("req-2", ARN, false, InvocationStatus.SUCCEEDED, null)); var spans = spanExporter.getFinishedSpanItems(); @@ -550,7 +642,18 @@ void allSpansShareTraceId_acrossInvocations() { 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)); + "op-1", + "step-1", + "STEP", + "Step", + null, + Instant.now(), + Instant.now(), + "SUCCEEDED", + null, + false, + null, + null)); plugin.onInvocationEnd(new InvocationEndInfo("req-1", ARN, true, InvocationStatus.PENDING, null)); var firstTraceId = spanExporter.getFinishedSpanItems().get(0).getTraceId(); spanExporter.reset(); @@ -579,6 +682,7 @@ void operationEnd_withoutStart_createsContinuationSpanWithLink() { "SUCCEEDED", null, false, + null, null)); plugin.onInvocationEnd(new InvocationEndInfo("req-2", ARN, false, InvocationStatus.SUCCEEDED, null)); @@ -650,7 +754,18 @@ void xrayExtraction_allSpansShareExtractedTraceId() { 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)); + "op-1", + "step-a", + "STEP", + "Step", + null, + Instant.now(), + Instant.now(), + "SUCCEEDED", + null, + false, + null, + null)); xrayPlugin.onInvocationEnd(new InvocationEndInfo("req-1", ARN, true, InvocationStatus.SUCCEEDED, null)); var spans = exporter.getFinishedSpanItems(); 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 018273260..38961b283 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 @@ -98,7 +98,18 @@ void defaultConstructor_usesGlobalSdkTracerProviderDirectly() { 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)); + "op-1", + "step", + "STEP", + "Step", + null, + Instant.now(), + Instant.now(), + "SUCCEEDED", + null, + false, + null, + null)); defaultPlugin.onInvocationEnd( new InvocationEndInfo("req-1", "arn:exec1", true, InvocationStatus.SUCCEEDED, null)); @@ -136,7 +147,18 @@ public ContextPropagators getPropagators() { 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)); + "op-1", + "step", + "STEP", + "Step", + null, + Instant.now(), + Instant.now(), + "SUCCEEDED", + null, + false, + null, + null)); defaultPlugin.onInvocationEnd( new InvocationEndInfo("req-1", "arn:exec1", true, InvocationStatus.SUCCEEDED, null)); @@ -260,6 +282,7 @@ void operationSpanName_usesOperationName_withoutPrefix() { "SUCCEEDED", null, false, + null, null)); plugin.onInvocationEnd(new InvocationEndInfo("req-1", "arn:exec1", true, InvocationStatus.SUCCEEDED, null)); @@ -299,7 +322,18 @@ void operationEnd_withAttempt_stampsAttemptNumberOnOperationSpan() { 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)); + "op-1", + "flaky", + "STEP", + "Step", + null, + Instant.now(), + Instant.now(), + "SUCCEEDED", + 3, + false, + null, + null)); plugin.onInvocationEnd(new InvocationEndInfo("req-1", "arn:exec1", true, InvocationStatus.SUCCEEDED, null)); var operationSpan = spanExporter.getFinishedSpanItems().stream() @@ -337,7 +371,18 @@ void operationEnd_withoutMatchingStart_stampsAttemptNumberOnContinuationSpan() { // 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)); + "op-1", + "flaky", + "STEP", + "Step", + null, + Instant.now(), + Instant.now(), + "SUCCEEDED", + 2, + false, + null, + null)); plugin.onInvocationEnd(new InvocationEndInfo("req-1", "arn:exec1", true, InvocationStatus.SUCCEEDED, null)); var continuationSpan = spanExporter.getFinishedSpanItems().stream() @@ -392,7 +437,7 @@ void operationStart_createsSpan_operationEnd_endsIt() { // Operation span ended at completion plugin.onOperationEnd(new OperationEndInfo( - "op-hash-1", "my-step", "STEP", "Step", null, start, end, "SUCCEEDED", null, false, null)); + "op-hash-1", "my-step", "STEP", "Step", null, start, end, "SUCCEEDED", null, false, null, null)); plugin.onInvocationEnd(new InvocationEndInfo("req-1", "arn:exec1", true, InvocationStatus.SUCCEEDED, null)); @@ -483,7 +528,18 @@ void operationEnd_withSuccess_setsOkOnOperationSpan() { 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)); + "op-1", + "step-ok", + "STEP", + "Step", + null, + Instant.now(), + Instant.now(), + "SUCCEEDED", + null, + false, + null, + null)); plugin.onInvocationEnd(new InvocationEndInfo("req-1", "arn:exec1", true, InvocationStatus.SUCCEEDED, null)); @@ -516,6 +572,7 @@ void operationEnd_withNonSuccessStatusAndNoError_leavesOperationSpanUnset() { "CANCELLED", null, false, + null, null)); plugin.onInvocationEnd(new InvocationEndInfo("req-1", "arn:exec1", true, InvocationStatus.SUCCEEDED, null)); @@ -544,6 +601,7 @@ void operationEnd_withoutMatchingStart_nonSuccessStatusAndNoError_leavesContinua "TIMED_OUT", null, false, + null, null)); plugin.onInvocationEnd(new InvocationEndInfo("req-1", "arn:exec1", true, InvocationStatus.SUCCEEDED, null)); @@ -564,7 +622,18 @@ void operationEnd_withNullStatusAndNoError_setsOkOnOperationSpan() { 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)); + "op-ctx", + "my-ctx", + "CONTEXT", + null, + null, + Instant.now(), + Instant.now(), + null, + null, + false, + null, + null)); plugin.onInvocationEnd(new InvocationEndInfo("req-1", "arn:exec1", true, InvocationStatus.SUCCEEDED, null)); @@ -588,7 +657,18 @@ void fullLifecycle_producesCorrectSpanHierarchy() { plugin.onUserFunctionEnd(new UserFunctionEndInfo( "op-1", "step-a", "STEP", "Step", null, Instant.now(), Instant.now(), false, 1, true, null)); plugin.onOperationEnd(new OperationEndInfo( - "op-1", "step-a", "STEP", "Step", null, Instant.now(), Instant.now(), "SUCCEEDED", null, false, null)); + "op-1", + "step-a", + "STEP", + "Step", + null, + Instant.now(), + Instant.now(), + "SUCCEEDED", + null, + false, + null, + null)); // Step 2: operation starts, user function runs, operation completes plugin.onOperationStart( @@ -598,7 +678,18 @@ void fullLifecycle_producesCorrectSpanHierarchy() { plugin.onUserFunctionEnd(new UserFunctionEndInfo( "op-2", "step-b", "STEP", "Step", null, Instant.now(), Instant.now(), false, 1, true, null)); plugin.onOperationEnd(new OperationEndInfo( - "op-2", "step-b", "STEP", "Step", null, Instant.now(), Instant.now(), "SUCCEEDED", null, false, null)); + "op-2", + "step-b", + "STEP", + "Step", + null, + Instant.now(), + Instant.now(), + "SUCCEEDED", + null, + false, + null, + null)); plugin.onInvocationEnd(new InvocationEndInfo("req-1", arn, true, InvocationStatus.SUCCEEDED, null)); @@ -684,7 +775,18 @@ void sampling_disabled_producesNoSpans() { sampledPlugin.onUserFunctionEnd(new UserFunctionEndInfo( "op-1", "step", "STEP", "Step", null, Instant.now(), Instant.now(), false, 1, true, null)); sampledPlugin.onOperationEnd(new OperationEndInfo( - "op-1", "step", "STEP", "Step", null, Instant.now(), Instant.now(), "SUCCEEDED", null, false, null)); + "op-1", + "step", + "STEP", + "Step", + null, + Instant.now(), + Instant.now(), + "SUCCEEDED", + null, + false, + null, + null)); sampledPlugin.onInvocationEnd( new InvocationEndInfo("req-1", "arn:exec1", true, InvocationStatus.SUCCEEDED, null)); @@ -731,7 +833,18 @@ void xrayExtraction_allSpansShareExtractedTraceId() { xrayPlugin.onUserFunctionEnd(new UserFunctionEndInfo( "op-1", "step-a", "STEP", "Step", null, Instant.now(), Instant.now(), false, 1, true, null)); xrayPlugin.onOperationEnd(new OperationEndInfo( - "op-1", "step-a", "STEP", "Step", null, Instant.now(), Instant.now(), "SUCCEEDED", null, false, null)); + "op-1", + "step-a", + "STEP", + "Step", + null, + Instant.now(), + Instant.now(), + "SUCCEEDED", + null, + false, + null, + null)); xrayPlugin.onInvocationEnd(new InvocationEndInfo("req-1", "arn:exec1", true, InvocationStatus.SUCCEEDED, null)); var spans = spanExporter.getFinishedSpanItems(); @@ -813,7 +926,18 @@ void xrayExtraction_multipleInvocations_sameTraceId_unifiedTrace() { 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)); + "op-1", + "step-1", + "STEP", + "Step", + null, + Instant.now(), + Instant.now(), + "SUCCEEDED", + null, + false, + null, + null)); xrayPlugin.onInvocationEnd(new InvocationEndInfo("req-1", "arn:exec1", true, InvocationStatus.PENDING, null)); // Second invocation (same execution, same X-Ray Root from backend) @@ -821,7 +945,18 @@ void xrayExtraction_multipleInvocations_sameTraceId_unifiedTrace() { 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)); + "op-2", + "step-2", + "STEP", + "Step", + null, + Instant.now(), + Instant.now(), + "SUCCEEDED", + null, + false, + null, + null)); xrayPlugin.onInvocationEnd( new InvocationEndInfo("req-2", "arn:exec1", false, InvocationStatus.SUCCEEDED, null)); @@ -898,6 +1033,7 @@ void operationEnd_withoutMatchingStart_createsContinuationSpanWithLink() { "SUCCEEDED", null, false, + null, null)); plugin.onInvocationEnd(new InvocationEndInfo("req-1", "arn:exec1", true, InvocationStatus.SUCCEEDED, null)); @@ -932,6 +1068,7 @@ void operationEnd_withoutMatchingStart_usesOperationStartTimestamp() { "SUCCEEDED", null, false, + null, null)); plugin.onInvocationEnd(new InvocationEndInfo("req-1", "arn:exec1", true, InvocationStatus.SUCCEEDED, null)); @@ -962,7 +1099,8 @@ void operationEnd_withoutMatchingStart_withError_setsErrorStatus() { "FAILED", null, false, - new RuntimeException("timed out"))); + new RuntimeException("timed out"), + null)); plugin.onInvocationEnd(new InvocationEndInfo("req-1", "arn:exec1", true, InvocationStatus.SUCCEEDED, null)); @@ -1054,6 +1192,7 @@ void childOperation_parentedToParentOperationSpan() { "SUCCEEDED", null, false, + null, null)); plugin.onOperationEnd(new OperationEndInfo( @@ -1067,6 +1206,7 @@ void childOperation_parentedToParentOperationSpan() { "SUCCEEDED", null, false, + null, null)); plugin.onInvocationEnd(new InvocationEndInfo("req-1", "arn:exec1", true, InvocationStatus.SUCCEEDED, null)); @@ -1103,7 +1243,18 @@ void multiInvocation_stepWaitStep_producesCorrectSpans() { plugin.onUserFunctionEnd(new UserFunctionEndInfo( "op-1", "step-A", "STEP", "Step", null, Instant.now(), Instant.now(), false, 1, true, null)); plugin.onOperationEnd(new OperationEndInfo( - "op-1", "step-A", "STEP", "Step", null, Instant.now(), Instant.now(), "SUCCEEDED", null, false, null)); + "op-1", + "step-A", + "STEP", + "Step", + null, + Instant.now(), + Instant.now(), + "SUCCEEDED", + null, + false, + null, + 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)); @@ -1117,7 +1268,18 @@ void multiInvocation_stepWaitStep_producesCorrectSpans() { // Invocation 2: wait completed between invocations, new step runs plugin.onInvocationStart(new InvocationInfo("req-2", arn, false, Instant.now())); plugin.onOperationEnd(new OperationEndInfo( - "op-2", "pause", "WAIT", "Wait", null, Instant.now(), Instant.now(), "SUCCEEDED", null, false, null)); + "op-2", + "pause", + "WAIT", + "Wait", + null, + Instant.now(), + Instant.now(), + "SUCCEEDED", + null, + false, + null, + null)); plugin.onOperationStart( new OperationInfo("op-3", "step-B", "STEP", "Step", null, Instant.now(), null, null, false)); plugin.onUserFunctionStart( @@ -1125,7 +1287,18 @@ void multiInvocation_stepWaitStep_producesCorrectSpans() { plugin.onUserFunctionEnd(new UserFunctionEndInfo( "op-3", "step-B", "STEP", "Step", null, Instant.now(), Instant.now(), false, 1, true, null)); plugin.onOperationEnd(new OperationEndInfo( - "op-3", "step-B", "STEP", "Step", null, Instant.now(), Instant.now(), "SUCCEEDED", null, false, null)); + "op-3", + "step-B", + "STEP", + "Step", + null, + Instant.now(), + Instant.now(), + "SUCCEEDED", + null, + false, + null, + null)); plugin.onInvocationEnd(new InvocationEndInfo("req-2", arn, false, InvocationStatus.SUCCEEDED, null)); var inv2Spans = spanExporter.getFinishedSpanItems(); @@ -1212,6 +1385,7 @@ void crossInvocation_stepRetry_attemptsParentedToRespectiveInvocations() { "SUCCEEDED", null, false, + null, null)); plugin.onInvocationEnd(new InvocationEndInfo("req-2", arn, false, InvocationStatus.SUCCEEDED, null)); @@ -1289,7 +1463,18 @@ void operationAndAttemptSpans_linkToWorkflowSpan() { plugin.onUserFunctionEnd(new UserFunctionEndInfo( "op-1", "step-a", "STEP", "Step", null, Instant.now(), Instant.now(), false, 1, true, null)); plugin.onOperationEnd(new OperationEndInfo( - "op-1", "step-a", "STEP", "Step", null, Instant.now(), Instant.now(), "SUCCEEDED", 1, false, null)); + "op-1", + "step-a", + "STEP", + "Step", + null, + Instant.now(), + Instant.now(), + "SUCCEEDED", + 1, + false, + null, + null)); plugin.onInvocationEnd(new InvocationEndInfo("req-1", "arn:exec-wf", true, InvocationStatus.SUCCEEDED, null)); var workflowId = spanByName("Workflow").getSpanId(); @@ -1315,7 +1500,18 @@ void operationLinksToWorkflow_withXRayContext() { 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)); + "op-1", + "step-a", + "STEP", + "Step", + null, + Instant.now(), + Instant.now(), + "SUCCEEDED", + null, + false, + null, + null)); xrayPlugin.onInvocationEnd(new InvocationEndInfo("req-1", "arn:exec1", true, InvocationStatus.SUCCEEDED, null)); var workflowId = exporter.getFinishedSpanItems().stream() diff --git a/sdk-integration-tests/src/test/java/software/amazon/lambda/durable/PluginIntegrationTest.java b/sdk-integration-tests/src/test/java/software/amazon/lambda/durable/PluginIntegrationTest.java index ad61c8953..34bead95b 100644 --- a/sdk-integration-tests/src/test/java/software/amazon/lambda/durable/PluginIntegrationTest.java +++ b/sdk-integration-tests/src/test/java/software/amazon/lambda/durable/PluginIntegrationTest.java @@ -333,6 +333,25 @@ void plugin_operationEnd_noError_whenOperationSucceeds() { assertNull(stepEnd.error(), "error should be null for successful step"); } + @Test + void plugin_operationEnd_includesResult_whenStepSucceeds() { + var plugin = new RecordingPlugin(); + var config = DurableConfig.builder().withPlugins(plugin).build(); + + var runner = LocalDurableTestRunner.create( + String.class, (input, context) -> context.step("my-step", String.class, stepCtx -> "task-a"), config); + + runner.runUntilComplete("input"); + + var stepEnd = plugin.operationEnds.stream() + .filter(info -> "my-step".equals(info.name())) + .findFirst() + .orElse(null); + assertNotNull(stepEnd, "onOperationEnd should fire for successful step"); + assertNotNull(stepEnd.result(), "result should be present for successful step"); + assertEquals("\"task-a\"", stepEnd.result(), "result should contain the serialized step output"); + } + @Test void plugin_operationStartAndEnd_balanced_forEmptyMap() { var plugin = new RecordingPlugin(); diff --git a/sdk/src/main/java/software/amazon/lambda/durable/plugin/OperationEndInfo.java b/sdk/src/main/java/software/amazon/lambda/durable/plugin/OperationEndInfo.java index 685d25476..1b49c2d48 100644 --- a/sdk/src/main/java/software/amazon/lambda/durable/plugin/OperationEndInfo.java +++ b/sdk/src/main/java/software/amazon/lambda/durable/plugin/OperationEndInfo.java @@ -18,6 +18,7 @@ * @param attempt the total number of attempts for retriable operations (STEP, WAIT_FOR_CONDITION) — null for others * @param isReplay true if this operation already existed in the execution state (completed in a prior invocation) * @param error non-null if the operation failed + * @param result the serialized result of the operation (may be null if the operation failed or has no result) * @deprecated This is a preview API that is experimental and may be changed or removed in future releases. */ @Deprecated @@ -32,4 +33,5 @@ public record OperationEndInfo( String status, Integer attempt, boolean isReplay, - Throwable error) {} + Throwable error, + String result) {} 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..110d90e71 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 @@ -69,7 +69,33 @@ public static OperationEndInfo toOperationEndInfo( ? operation.stepDetails().attempt() : null, isReplay, - error); + error, + extractResult(operation)); + } + + /** + * Extracts the serialized result from an operation based on its type. Returns null if the operation has no result + * (e.g., failed or still running). + */ + private static String extractResult(Operation operation) { + if (operation == null || operation.type() == null) { + return null; + } + return switch (operation.type()) { + case STEP -> + operation.stepDetails() != null ? operation.stepDetails().result() : null; + case CHAINED_INVOKE -> + operation.chainedInvokeDetails() != null + ? operation.chainedInvokeDetails().result() + : null; + case CALLBACK -> + operation.callbackDetails() != null + ? operation.callbackDetails().result() + : null; + case CONTEXT -> + operation.contextDetails() != null ? operation.contextDetails().result() : null; + default -> null; + }; } /** 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..65e5c078a 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 @@ -8,6 +8,8 @@ import org.junit.jupiter.api.Test; import software.amazon.awssdk.services.lambda.model.Operation; import software.amazon.awssdk.services.lambda.model.OperationStatus; +import software.amazon.awssdk.services.lambda.model.OperationType; +import software.amazon.awssdk.services.lambda.model.StepDetails; import software.amazon.lambda.durable.model.OperationIdentifier; import software.amazon.lambda.durable.model.OperationSubType; @@ -96,6 +98,35 @@ void toOperationEndInfo_nullError_forSuccess() { assertNull(info.error()); } + @Test + void toOperationEndInfo_extractsResult_fromSucceededStep() { + var operation = Operation.builder() + .startTimestamp(START) + .endTimestamp(END) + .type(OperationType.STEP) + .status(OperationStatus.SUCCEEDED) + .stepDetails(StepDetails.builder().result("\"hello\"").build()) + .build(); + + var info = PluginInfoConverter.toOperationEndInfo(operation, STEP_IDENTIFIER, null, false, null); + + assertEquals("\"hello\"", info.result()); + } + + @Test + void toOperationEndInfo_resultIsNull_whenOperationHasNoResult() { + var operation = Operation.builder() + .startTimestamp(START) + .endTimestamp(END) + .type(OperationType.WAIT) + .status(OperationStatus.SUCCEEDED) + .build(); + + var info = PluginInfoConverter.toOperationEndInfo(operation, WAIT_IDENTIFIER, null, false, null); + + assertNull(info.result()); + } + // ─── toUserFunctionStartInfo ──────────────────────────────────────── @Test 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..65588bf4f 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 @@ -151,7 +151,18 @@ private static OperationInfo operationInfo() { private static OperationEndInfo operationEndInfo() { return new OperationEndInfo( - "op-1", "test-step", "STEP", null, null, Instant.now(), Instant.now(), "SUCCEEDED", null, false, null); + "op-1", + "test-step", + "STEP", + null, + null, + Instant.now(), + Instant.now(), + "SUCCEEDED", + null, + false, + null, + null); } private static OperationChangeInfo operationChangeInfo() {