From fdbfe492ab6da2298ce7d5fad10435a7a1ba1887 Mon Sep 17 00:00:00 2001 From: Ayushi Ahjolia Date: Wed, 5 Aug 2026 16:35:06 -0700 Subject: [PATCH 1/2] 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() { From 36eae1a7e8be6737132621d30b6843625b61aee8 Mon Sep 17 00:00:00 2001 From: Alex Wang Date: Thu, 13 Aug 2026 19:24:42 +0000 Subject: [PATCH 2/2] fix(plugin): guard result on success, keep it out of toString, keep the old arity Review feedback on the result accessor: - Only a SUCCEEDED operation reports a result. A wait-for-condition is checkpointed as STEP and reuses stepDetails().result() to carry its intermediate check-loop state between attempts, so a failed one could surface that state as its result, contradicting the documented null-on-failure contract. Verified: without the guard the new test reports {"polls":2} for a FAILED operation. - Override toString() to omit result. Records render every component, so plugins that log the info object whole would have begun emitting customer payloads (possibly secrets or personal data). Output is otherwise unchanged; matches the invocation records from #622. - Add a constructor at the previous 11-argument arity delegating with a null result, so existing callers keep compiling and linking. - Cover the CHAINED_INVOKE, CALLBACK, and CONTEXT extraction branches, which had no direct test. --- .../durable/plugin/OperationEndInfo.java | 50 ++++++++- .../durable/plugin/PluginInfoConverter.java | 8 +- .../plugin/PluginInfoConverterTest.java | 100 ++++++++++++++++++ 3 files changed, 156 insertions(+), 2 deletions(-) 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 e1133e917..b61df330e 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 @@ -34,4 +34,52 @@ public record OperationEndInfo( Integer attempt, boolean isReplay, @Experimental Throwable error, - @Experimental String result) {} + @Experimental String result) { + + /** + * Creates operation-end information without a result. + * + *

Retained so callers written before {@code result} was added keep compiling and linking; {@code result} + * resolves to null. + * + * @param id operation ID + * @param name human-readable operation name (may be null) + * @param type operation type + * @param subType operation sub-type (may be null) + * @param parentId parent operation ID (null for root-level operations) + * @param startTimestamp when the operation started + * @param endTimestamp when the operation ended + * @param status the operation's terminal status + * @param attempt the total number of attempts for retriable operations + * @param isReplay true if this operation already existed in the execution state + * @param error non-null if the operation failed + */ + public OperationEndInfo( + String id, + String name, + String type, + String subType, + String parentId, + Instant startTimestamp, + Instant endTimestamp, + String status, + Integer attempt, + boolean isReplay, + Throwable error) { + this(id, name, type, subType, parentId, startTimestamp, endTimestamp, status, attempt, isReplay, error, null); + } + + /** + * Returns a representation that omits {@code result}. + * + *

The generated representation would render the operation result, so plugins that log this object whole would + * start emitting customer payloads, potentially including secrets or personal data. Read the component explicitly + * to record it. + */ + @Override + public String toString() { + return "OperationEndInfo[id=" + id + ", name=" + name + ", type=" + type + ", subType=" + subType + + ", parentId=" + parentId + ", startTimestamp=" + startTimestamp + ", endTimestamp=" + endTimestamp + + ", status=" + status + ", attempt=" + attempt + ", isReplay=" + isReplay + ", error=" + 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 ad91bfe18..59b9e4a3a 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 @@ -6,6 +6,7 @@ import java.util.Collection; import java.util.stream.Collectors; import software.amazon.awssdk.services.lambda.model.Operation; +import software.amazon.awssdk.services.lambda.model.OperationStatus; import software.amazon.lambda.durable.model.OperationIdentifier; import software.amazon.lambda.durable.operation.BaseDurableOperation; @@ -71,9 +72,14 @@ public static OperationEndInfo toOperationEndInfo( /** * 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). + * + *

Only a SUCCEEDED operation reports a result. The status guard is required, not just defensive: a + * wait-for-condition is checkpointed as {@link OperationType#STEP} and reuses {@code stepDetails().result()} to + * carry its intermediate check-loop state between attempts, so a failed one can still hold state that must not be + * surfaced as that operation's result. */ private static String extractResult(Operation operation) { - if (operation == null || operation.type() == null) { + if (operation == null || operation.type() == null || operation.status() != OperationStatus.SUCCEEDED) { return null; } return switch (operation.type()) { 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 65e5c078a..b4686a16b 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 @@ -6,6 +6,9 @@ import java.time.Instant; import org.junit.jupiter.api.Test; +import software.amazon.awssdk.services.lambda.model.CallbackDetails; +import software.amazon.awssdk.services.lambda.model.ChainedInvokeDetails; +import software.amazon.awssdk.services.lambda.model.ContextDetails; import software.amazon.awssdk.services.lambda.model.Operation; import software.amazon.awssdk.services.lambda.model.OperationStatus; import software.amazon.awssdk.services.lambda.model.OperationType; @@ -113,6 +116,103 @@ void toOperationEndInfo_extractsResult_fromSucceededStep() { assertEquals("\"hello\"", info.result()); } + @Test + void toOperationEndInfo_extractsResult_fromSucceededChainedInvoke() { + var operation = Operation.builder() + .startTimestamp(START) + .endTimestamp(END) + .type(OperationType.CHAINED_INVOKE) + .status(OperationStatus.SUCCEEDED) + .chainedInvokeDetails( + ChainedInvokeDetails.builder().result("\"invoked\"").build()) + .build(); + + var info = PluginInfoConverter.toOperationEndInfo(operation, STEP_IDENTIFIER, null, false, null); + + assertEquals("\"invoked\"", info.result()); + } + + @Test + void toOperationEndInfo_extractsResult_fromSucceededCallback() { + var operation = Operation.builder() + .startTimestamp(START) + .endTimestamp(END) + .type(OperationType.CALLBACK) + .status(OperationStatus.SUCCEEDED) + .callbackDetails( + CallbackDetails.builder().result("\"called-back\"").build()) + .build(); + + var info = PluginInfoConverter.toOperationEndInfo(operation, STEP_IDENTIFIER, null, false, null); + + assertEquals("\"called-back\"", info.result()); + } + + @Test + void toOperationEndInfo_extractsResult_fromSucceededContext() { + var operation = Operation.builder() + .startTimestamp(START) + .endTimestamp(END) + .type(OperationType.CONTEXT) + .status(OperationStatus.SUCCEEDED) + .contextDetails( + ContextDetails.builder().result("\"child-done\"").build()) + .build(); + + var info = PluginInfoConverter.toOperationEndInfo(operation, STEP_IDENTIFIER, null, false, null); + + assertEquals("\"child-done\"", info.result()); + } + + @Test + void toOperationEndInfo_resultIsNull_whenFailedWaitForConditionRetainsCheckpointState() { + // A wait-for-condition is checkpointed as STEP and reuses stepDetails().result() to carry its + // intermediate check-loop state between attempts, so a failed one can still hold state. + var operation = Operation.builder() + .startTimestamp(START) + .endTimestamp(END) + .type(OperationType.STEP) + .status(OperationStatus.FAILED) + .stepDetails( + StepDetails.builder().attempt(3).result("{\"polls\":2}").build()) + .build(); + + var info = PluginInfoConverter.toOperationEndInfo(operation, STEP_IDENTIFIER, null, false, null); + + assertNull(info.result(), "a failed operation must not report intermediate state as its result"); + } + + @Test + void operationEndInfo_compatibilityConstructor_leavesResultNull() { + var info = new OperationEndInfo( + OPERATION_ID, OPERATION_NAME, "STEP", "Step", PARENT_ID, START, END, "SUCCEEDED", 1, false, null); + + assertNull(info.result()); + } + + @Test + void operationEndInfo_toString_omitsResult() { + var info = new OperationEndInfo( + OPERATION_ID, + OPERATION_NAME, + "STEP", + "Step", + PARENT_ID, + START, + END, + "SUCCEEDED", + 1, + false, + null, + "s3cret-result"); + + var rendered = info.toString(); + + assertFalse(rendered.contains("s3cret-result"), "operation result must not leak into logs"); + assertTrue(rendered.contains(OPERATION_ID)); + assertTrue(rendered.contains("SUCCEEDED")); + } + @Test void toOperationEndInfo_resultIsNull_whenOperationHasNoResult() { var operation = Operation.builder()