Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
56 commits
Select commit Hold shift + click to select a range
527d8be
fix(otel): scope root handler instrumentation to its execution
Oct 3, 2026
06727e7
fix(otel): finalize handler scopes before invocation response
Oct 3, 2026
ba840d2
fix(plugin): preserve wrapped fatal handler-scope failures
Oct 3, 2026
6269b00
fix(ci): keep PR review intake read-only
Oct 3, 2026
3e76962
fix(execution): bound handler scope cleanup and preserve control flow
Oct 3, 2026
8132a08
test: verify earlier context cleanup after a scope fatal
Oct 3, 2026
42dc7e1
fix: preserve plugin ABI and observe scope fatals after shutdown
Oct 3, 2026
5b2d85f
fix: isolate handler scope openers from legacy subclass methods
Oct 3, 2026
3d94dd4
fix: preserve caller and handler MDC around optional plugins
zhongkechen Oct 6, 2026
4d81d8c
test: cover fatal cleanup during invocation end dispatch
zhongkechen Oct 6, 2026
ddcd77b
fix: preserve invocation end completion-thread dispatch
zhongkechen Oct 6, 2026
de9ea32
ci: queue serial conformance runs without replacing pending work
zhongkechen Oct 6, 2026
07595bb
Merge branch 'main' into fix/otel-handler-context
zhongkechen Oct 6, 2026
1f7ddb8
Merge main after factory migration revert
Oct 6, 2026
b114ab7
Register the replay validation test caller as active
Oct 6, 2026
ee73b45
ci: run OTel conformance on validated CodeBuild Java 21
Oct 7, 2026
0e0daeb
fix: preserve end-hook completion ownership and MDC lifetime
Oct 7, 2026
b5c688c
fix: settle handler future when MDC capture fails
Oct 7, 2026
d5bd3b8
fix: wake invocation caller on reported scope fatal
zhongkechen Oct 7, 2026
768fd0f
Merge main OpenTelemetry API compatibility fix
Oct 7, 2026
b1f46e2
Merge main exclusive OpenTelemetry view validation
Oct 7, 2026
0d8d860
fix: preserve execution outcomes across end MDC failures
Oct 7, 2026
5a7f549
fix: accept concrete handler scope return types
Oct 7, 2026
a9d3285
fix: contain malformed handler scope diagnostics
Oct 7, 2026
923a914
test: return one checkpoint state per operation in unit mocks
zhongkechen Oct 7, 2026
e72b3f2
test: include failed execution details in executor reuse checks
Oct 7, 2026
7c52725
Merge branch 'main' into fix/otel-handler-context
zhongkechen Oct 7, 2026
fac8651
refactor: finalize invocation hooks on the handler thread
zhongkechen Oct 7, 2026
c1a3310
ci: pin OTel queue workflow to published squash commit
Oct 7, 2026
8c7d121
fix(plugin): enforce invocation lifecycle boundaries
Oct 7, 2026
cacd730
fix: propagate fatal MDC initialization failures before startup
Oct 8, 2026
b1fb4eb
fix: preserve invocation end failures through scope cleanup
Oct 8, 2026
e9a07b2
test: assert stored parallel tolerance outcome across racing branches
Oct 8, 2026
cdb5c0c
fix: preserve delivery and settled invocation outcomes
Oct 8, 2026
1dc8136
docs: define the SDK invocation-end response boundary
Oct 8, 2026
7c09c60
fix: retain ready condition work across checkpoint responses
Oct 8, 2026
7a4ff80
fix: inspect MDC initialization causes without cycling
Oct 8, 2026
58ef569
fix: preserve terminal wait-for-condition outcomes on replay
Oct 8, 2026
565d696
ci: retry readonly layer lookup for the existing 20 OTel cases
Oct 8, 2026
bcb9a09
fix: resume condition polling outside the checkpoint batcher
Oct 8, 2026
4b048d9
fix: yield READY checks to queued user operations
Oct 8, 2026
23abc38
fix: consume durable sampling intent before processor callbacks
Oct 8, 2026
85a09eb
fix: drain owned condition continuations before manager shutdown
Oct 8, 2026
c0d7e06
fix: retain live operation parents for native span clocks
Oct 8, 2026
c842f5a
fix: preserve preparation failures and isolate sampling carriers
Oct 8, 2026
478c183
fix: preserve foreign sampling metadata and execution isolation
Oct 8, 2026
c96e11a
fix: select manager control flow before waking handlers
Oct 8, 2026
5f23e7d
fix: retain visible replacement sampler root policy
Oct 8, 2026
80110fc
fix: propagate coordinator and post-End fatal failures
Oct 8, 2026
007031a
fix: preserve owned continuation failures and rejected admission
Oct 8, 2026
6794136
fix: retain retry diagnostics and fatal identity through End cleanup
Oct 8, 2026
78e3b48
fix: clear inheritable MDC after ordinary restoration failure
Oct 8, 2026
43b8c34
fix: coordinate step retries outside the checkpoint batcher
Oct 9, 2026
0e55c1a
fix: retain fatal cleanup and preparation error diagnostics
Oct 9, 2026
d148379
fix: select continuation failures atomically with close
Oct 9, 2026
be44fd7
fix: prune completed polls before checkpoint batching
Oct 9, 2026
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 6 additions & 1 deletion .github/scripts/verify_otel_api_compatibility.py
Original file line number Diff line number Diff line change
Expand Up @@ -142,11 +142,14 @@ def run_matrix(args: argparse.Namespace) -> int:
for view in ("otel-invocation", "otel-execution"):
name = f"{label}-api{version}-{view}"
negative = label == "old-old" and version == "1.49.0"
incompatible_core = label == "old-new"
case: dict[str, object] = {"name": name, "expected_negative_control": negative,
"expected_core_rejection": incompatible_core,
"api": jar_facts(api), "context": jar_facts(context)}
command = [args.java, "-cp", os.pathsep.join(map(str, [classes, core, plugin, *dependencies])),
PROBE, str(core), str(plugin), str(api), str(context), view,
str(negative).lower(), str(version == "1.66.0").lower()]
str(negative).lower(), str(version == "1.66.0").lower(),
str(incompatible_core).lower()]
try:
log = output / f"{name}.log"
execute(command, log, env=probe_environment(view), timeout=90)
Expand All @@ -155,6 +158,8 @@ def run_matrix(args: argparse.Namespace) -> int:
raise RuntimeError("Probe did not report successful completion")
if negative and "NEGATIVE_CONTROL_REPRODUCED" not in contents:
raise RuntimeError("Released negative control did not reproduce the reported failure")
if incompatible_core and "CORE_LIFECYCLE_REJECTION_CONFIRMED" not in contents:
raise RuntimeError("Old core must reject the new plugin before invocation startup")
case["passed"] = True
except RuntimeError as error:
failures += 1
Expand Down
3 changes: 2 additions & 1 deletion .github/workflows/otel-conformance-tests.yml
Original file line number Diff line number Diff line change
Expand Up @@ -68,8 +68,9 @@ jobs:
actions: write
contents: read
id-token: write
uses: aws/aws-durable-execution-conformance-tests/.github/workflows/opentelemetry-orchestrator.yml@a66037abbbfa55fde97f714e30f0bc262edefd63
uses: aws/aws-durable-execution-conformance-tests/.github/workflows/opentelemetry-orchestrator.yml@1768f32d61d958b974ce9e66748df72e23bc19b2
with:
case_count: 20
runs_on: ${{ github.event_name == 'pull_request' && github.event.pull_request.head.repo.full_name != github.repository && 'ubuntu-latest' || format('codebuild-github-actions-runner-{0}-{1}', github.run_id, github.run_attempt) }}
language: java
resource_prefix: j
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@
import software.amazon.awssdk.regions.Region;
import software.amazon.awssdk.services.lambda.LambdaClient;
import software.amazon.awssdk.services.lambda.model.ErrorObject;
import software.amazon.awssdk.services.lambda.model.Event;
import software.amazon.awssdk.services.lambda.model.OperationStatus;
import software.amazon.awssdk.services.sts.StsClient;
import software.amazon.lambda.durable.TypeToken;
Expand Down Expand Up @@ -818,7 +819,20 @@ void testConcurrentWaitForConditionExample() {
lambdaClient);
var result = runner.run(new ConcurrentWaitForConditionExample.Input(3, 100, 50));

assertEquals(ExecutionStatus.SUCCEEDED, result.getStatus());
try {
assertEquals(ExecutionStatus.SUCCEEDED, result.getStatus());
} catch (RuntimeException | AssertionError failure) {
try {
// Preserve the already-fetched service timeline without payloads, callback IDs, or checkpoint tokens.
var history = result.getHistoryEvents().stream()
.map(CloudBasedIntegrationTest::historyMetadata)
.toList();
System.err.println("Concurrent wait-for-condition history: " + new JacksonSerDes().serialize(history));
} catch (RuntimeException diagnosticsFailure) {
failure.addSuppressed(diagnosticsFailure);
}
throw failure;
}

// Verify each operation finished with 3 attempts
var allOperationsOutput = result.getResult();
Expand All @@ -840,6 +854,30 @@ void testConcurrentWaitForConditionExample() {
}
}

private static Map<String, Object> historyMetadata(Event event) {
var row = new HashMap<String, Object>();
row.put("eventId", event.eventId());
row.put("type", event.eventTypeAsString());
row.put("operationId", event.id());
row.put("parentId", event.parentId());
row.put("name", event.name());
row.put("at", String.valueOf(event.eventTimestamp()));
var retries = event.stepSucceededDetails() != null
? event.stepSucceededDetails().retryDetails()
: event.stepFailedDetails() != null ? event.stepFailedDetails().retryDetails() : null;
if (retries != null) {
row.put("attempt", retries.currentAttempt());
row.put("nextAttemptDelaySeconds", retries.nextAttemptDelaySeconds());
}
var invocation = event.invocationCompletedDetails();
if (invocation != null) {
row.put("invocationStart", String.valueOf(invocation.startTimestamp()));
row.put("invocationEnd", String.valueOf(invocation.endTimestamp()));
row.put("invocationFailed", invocation.error() != null);
}
return row;
}

@Test
void testPluginExample() {
var runner =
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -5,8 +5,27 @@
import static org.junit.jupiter.api.Assertions.*;

import java.util.List;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.Executors;
import java.util.concurrent.TimeUnit;
import java.util.stream.Collectors;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.params.ParameterizedTest;
import org.junit.jupiter.params.provider.ValueSource;
import software.amazon.awssdk.services.lambda.model.OperationStatus;
import software.amazon.awssdk.services.lambda.model.OperationType;
import software.amazon.lambda.durable.DurableConfig;
import software.amazon.lambda.durable.TypeToken;
import software.amazon.lambda.durable.model.ConcurrencyCompletionStatus;
import software.amazon.lambda.durable.model.ExecutionStatus;
import software.amazon.lambda.durable.model.ParallelResult;
import software.amazon.lambda.durable.plugin.DurableExecutionPlugin;
import software.amazon.lambda.durable.plugin.OperationEndInfo;
import software.amazon.lambda.durable.plugin.UserFunctionEndInfo;
import software.amazon.lambda.durable.plugin.UserFunctionStartInfo;
import software.amazon.lambda.durable.serde.JacksonSerDes;
import software.amazon.lambda.durable.testing.LocalDurableTestRunner;

class ParallelFailureToleranceExampleTest {
Expand Down Expand Up @@ -41,19 +60,104 @@ void succeedsWhenAllBranchesSucceed() {
assertEquals(3, output.succeeded());
}

@Test
void failsWhenFailuresExceedTolerance() {
@ParameterizedTest
@ValueSource(booleans = {false, true})
void failsWhenFailuresExceedTolerance(boolean holdHealthyBranch) throws Exception {
var handler = new ParallelFailureToleranceExample();
var runner = LocalDurableTestRunner.create(ParallelFailureToleranceExample.Input.class, handler);
var probe = new CompletionProbe(holdHealthyBranch);
var config = DurableConfig.builder().withPlugins(probe.newPlugin()).build();
var runner = LocalDurableTestRunner.create(ParallelFailureToleranceExample.Input.class, handler)
.withDurableConfig(config);
var caller = Executors.newSingleThreadExecutor();
try {
var input = new ParallelFailureToleranceExample.Input(List.of("svc-a", "bad-svc-b", "bad-svc-c"), 1, 2);
var invocation = caller.submit(() -> runner.runUntilComplete(input));
if (holdHealthyBranch) {
await(probe.parallelStored);
probe.releaseHealthy.countDown();
}
var result = invocation.get(5, TimeUnit.SECONDS);
assertEquals(ExecutionStatus.SUCCEEDED, result.getStatus());
var output = result.getResult(ParallelFailureToleranceExample.Output.class);
assertEquals(2, output.failed());

// 2 bad services, toleratedFailureCount=1 — second failure exceeds tolerance
var input = new ParallelFailureToleranceExample.Input(List.of("svc-a", "bad-svc-b", "bad-svc-c"), 1, 2);
var result = runner.runUntilComplete(input);
var operation = result.getOperation("call-services");
assertEquals(OperationStatus.SUCCEEDED, operation.getStatus());
var stored = new JacksonSerDes()
.deserialize(operation.getContextDetails().result(), TypeToken.get(ParallelResult.class));
assertEquals(ConcurrencyCompletionStatus.FAILURE_TOLERANCE_EXCEEDED, stored.completionStatus());
assertFalse(stored.completionStatus().isSucceeded());
assertEquals(3, stored.size());
assertEquals(
List.of(ParallelResult.Status.FAILED, ParallelResult.Status.FAILED),
stored.statuses().subList(1, 3));
assertTrue(List.of(ParallelResult.Status.SUCCEEDED, ParallelResult.Status.SKIPPED)
.contains(stored.statuses().get(0)));
assertEquals(output.succeeded(), stored.succeeded());
assertEquals(output.failed(), stored.failed());
assertEquals(1 - output.succeeded(), stored.skipped());
if (holdHealthyBranch) assertEquals(0, output.succeeded());

assertEquals(ExecutionStatus.SUCCEEDED, result.getStatus());
var completedCalls = result.getOperations().stream()
.filter(op -> op.getType() == OperationType.STEP && op.isCompleted())
.collect(Collectors.toMap(op -> op.getName(), op -> probe.completedCalls.get(op.getName())));
assertEquals(1, completedCalls.get("invoke-bad-svc-b"));
assertEquals(1, completedCalls.get("invoke-bad-svc-c"));
var replay = runner.run(input);
assertEquals(ExecutionStatus.SUCCEEDED, replay.getStatus());
assertEquals(output, replay.getResult(ParallelFailureToleranceExample.Output.class));
completedCalls.forEach((name, count) -> assertEquals(
count,
probe.completedCalls.get(name),
"Completed step bodies must not run again on replay: " + name));
} finally {
probe.releaseHealthy.countDown();
caller.shutdownNow();
assertTrue(caller.awaitTermination(5, TimeUnit.SECONDS));
}
}

var output = result.getResult(ParallelFailureToleranceExample.Output.class);
assertEquals(2, output.failed());
assertEquals(1, output.succeeded());
private static void await(CountDownLatch latch) {
try {
assertTrue(latch.await(5, TimeUnit.SECONDS), "Controlled branch scheduling was not released");
} catch (InterruptedException failure) {
Thread.currentThread().interrupt();
throw new AssertionError(failure);
}
}

private static final class CompletionProbe {
final boolean holdHealthy;
final CountDownLatch healthyStarted = new CountDownLatch(1);
final CountDownLatch releaseHealthy = new CountDownLatch(1);
final CountDownLatch parallelStored = new CountDownLatch(1);
final Map<String, Integer> completedCalls = new ConcurrentHashMap<>();

CompletionProbe(boolean holdHealthy) {
this.holdHealthy = holdHealthy;
}

DurableExecutionPlugin newPlugin() {
return new DurableExecutionPlugin() {
@Override
public void onUserFunctionStart(UserFunctionStartInfo info) {
if (!holdHealthy || !"STEP".equals(info.type())) return;
if ("invoke-svc-a".equals(info.name())) {
healthyStarted.countDown();
await(releaseHealthy);
} else if (info.name().startsWith("invoke-bad-")) await(healthyStarted);
}

@Override
public void onUserFunctionEnd(UserFunctionEndInfo info) {
if ("STEP".equals(info.type())) completedCalls.merge(info.name(), 1, Integer::sum);
}

@Override
public void onOperationEnd(OperationEndInfo info) {
if ("call-services".equals(info.name())) parallelStored.countDown();
}
};
}
}
}
Loading
Loading