Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
25 commits
Select commit Hold shift + click to select a range
b0a33bc
fix: isolate incompatible OpenTelemetry API linkage failures
zhongkechen Oct 6, 2026
b66e520
Revert "Revert "refactor(plugin)!: create one plugin instance per inv…
zhongkechen Oct 6, 2026
d5ef07f
Merge branch 'main' into fix/otel-api-linkage-compat-763
zhongkechen Oct 6, 2026
c0ec3ae
refactor: inherit OTel linkage fix from PR 780
zhongkechen Oct 6, 2026
10e8f11
test: verify installed OpenTelemetry API compatibility
zhongkechen Oct 6, 2026
ebe0fcd
test: adapt inherited API compatibility probe to 3.x factories
zhongkechen Oct 6, 2026
3b2d6f6
ci: run OTel conformance on validated CodeBuild Java 21
Oct 7, 2026
3a49c5a
ci: run OTel conformance on validated CodeBuild Java 21
Oct 7, 2026
c8b48be
Merge linkage base CI update while retaining 24-case factory tests
Oct 7, 2026
7c50701
ci: trigger OTel compatibility checks for module and harness changes
Oct 7, 2026
f4e4e1d
ci: use read-only resolver for review-comment intake
Oct 7, 2026
8f9a53b
Merge updated compatibility CI base into held factory branch
Oct 7, 2026
987a4e8
fix: wake fatal scope observers without ending owner cleanup
zhongkechen Oct 7, 2026
9333fd4
fix: bound fatal peer cleanup and contain malformed failure chains
zhongkechen Oct 7, 2026
58e3740
fix: preserve operation-owner cleanup during fatal cancellation
Oct 7, 2026
78373ab
fix: reject user function admission after a peer plugin fatal
Oct 7, 2026
5cf04d1
fix: retain owner cleanup through bounded plugin failures
Oct 7, 2026
89d96c3
Merge main after squash of already-integrated PR #780
Oct 7, 2026
ff2709b
Merge main exclusive OpenTelemetry view validation
Oct 7, 2026
95fc394
fix: accept concrete handler scope return types
Oct 7, 2026
6e512e7
fix: preserve execution outcomes across end MDC failures
Oct 7, 2026
2531949
fix: report worker MDC fatals before end dispatch
Oct 7, 2026
7b19a27
test: verify parallel tolerance without racing completion assumptions
Oct 7, 2026
fc75e95
fix: bound peer cleanup for selected invocation fatals
Oct 7, 2026
4045921
Merge main CI runner updates into PR branch
Oct 7, 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
79 changes: 36 additions & 43 deletions .github/scripts/verify_otel_api_compatibility.py
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
#!/usr/bin/env python3
# Copyright Amazon.com, Inc. or its affiliates. All Rights Reserved.
# SPDX-License-Identifier: Apache-2.0
"""Exercise real released/candidate core and plugin artifacts with two visible OTel APIs.
"""Exercise the current 3.x core, plugin, and testing artifacts with two visible OTel APIs.

Dependency resolution and every probe are required: a network, compilation, or case
failure returns nonzero. No production dependency versions are modified.
Expand Down Expand Up @@ -105,67 +105,59 @@ def run_matrix(args: argparse.Namespace) -> int:
output.mkdir(parents=True, exist_ok=True)
fixture = root / "otel-plugin/src/test/compatibility/b1"
cp = resolve_classpaths(fixture, output, args.maven)
released_core = artifact(cp["1.66.0"], CORE, RELEASED_VERSION)
released_plugin = artifact(cp["1.66.0"], PLUGIN, RELEASED_VERSION)
new_core = args.new_core.resolve() if args.new_core else candidate_jar(root, "sdk", CORE)
new_plugin = args.new_plugin.resolve() if args.new_plugin else candidate_jar(root, "otel-plugin", PLUGIN)
for jar in (new_core, new_plugin):
if not jar.is_file():
raise RuntimeError(f"Candidate artifact missing: {jar}")
candidate_inputs = {"core": jar_facts(new_core), "plugin": jar_facts(new_plugin)}
new_core = snapshot_candidate(new_core, output)
new_plugin = snapshot_candidate(new_plugin, output)
testing_name = CORE + "-testing"
inputs = {
"core": args.new_core.resolve() if args.new_core else candidate_jar(root, "sdk", CORE),
"plugin": args.new_plugin.resolve() if args.new_plugin else candidate_jar(root, "otel-plugin", PLUGIN),
"testing": args.new_testing.resolve() if args.new_testing else candidate_jar(root, "sdk-testing", testing_name),
}
candidate_inputs = {name: jar_facts(path) for name, path in inputs.items()}
selected = {name: snapshot_candidate(path, output) for name, path in inputs.items()}
excluded = {f"{name}-{RELEASED_VERSION}.jar" for name in (CORE, PLUGIN, testing_name)}
dependencies = {version: [path for path in paths if path.name not in excluded]
for version, paths in cp.items()}
classes = output / "classes"
classes.mkdir(exist_ok=True)
execute([args.javac, "--release", "17", "-classpath", os.pathsep.join(map(str, cp["1.66.0"])),
compile_cp = [*selected.values(), *dependencies["1.66.0"]]
execute([args.javac, "--release", "17", "-classpath", os.pathsep.join(map(str, compile_cp)),
"-d", str(classes), str(fixture / "InstalledApiProbe.java")], output / "compile.log")
services = classes / "META-INF/services"
services.mkdir(parents=True, exist_ok=True)
(services / "software.amazon.lambda.durable.plugin.DurableExecutionPluginProvider").write_text(
PROBE + "$HealthyProvider\n")
report: dict[str, object] = {
"released_core": jar_facts(released_core), "released_plugin": jar_facts(released_plugin),
"new_core": jar_facts(new_core), "new_plugin": jar_facts(new_plugin),
"contract": "Current 3.x factory API; cross-major core/plugin mixtures are unsupported.",
"candidate_inputs": candidate_inputs,
"candidate_snapshots": {name: jar_facts(path) for name, path in selected.items()},
"cases": [], "agent_coverage": "This matrix is visible-API skew, not a deployed Java-agent test.",
}
cases: list[dict[str, object]] = report["cases"] # type: ignore[assignment]
failures = 0
pairs = {"old-old": (released_core, released_plugin), "new-old": (new_core, released_plugin),
"old-new": (released_core, new_plugin), "new-new": (new_core, new_plugin)}
for version in API_VERSIONS:
api = artifact(cp[version], "opentelemetry-api", version)
context = artifact(cp[version], "opentelemetry-context", version)
dependencies = [p for p in cp[version] if p.name not in
(f"{CORE}-{RELEASED_VERSION}.jar", f"{PLUGIN}-{RELEASED_VERSION}.jar")]
for label, (core, plugin) in pairs.items():
for view in ("otel-invocation", "otel-execution"):
name = f"{label}-api{version}-{view}"
negative = label == "old-old" and version == "1.49.0"
case: dict[str, object] = {"name": name, "expected_negative_control": negative,
"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()]
try:
log = output / f"{name}.log"
execute(command, log, env=probe_environment(view), timeout=90)
contents = log.read_text(errors="replace")
if "COMPAT_PASS " not in contents:
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")
case["passed"] = True
except RuntimeError as error:
failures += 1
case.update(passed=False, error=str(error))
cases.append(case)
(output / "results.json").write_text(json.dumps(report, indent=2) + "\n")
print(f"{'PASS' if case['passed'] else 'FAIL'} {name}", flush=True)
for view in ("otel-invocation", "otel-execution"):
name = f"current3x-api{version}-{view}"
case: dict[str, object] = {"name": name, "api": jar_facts(api), "context": jar_facts(context)}
command = [args.java, "-cp", os.pathsep.join(map(str, [classes, *selected.values(), *dependencies[version]])),
PROBE, str(selected["core"]), str(selected["plugin"]), str(api), str(context), view,
str(version == "1.66.0").lower(), str(selected["testing"])]
try:
log = output / f"{name}.log"
execute(command, log, env=probe_environment(view), timeout=90)
if "COMPAT_PASS " not in log.read_text(errors="replace"):
raise RuntimeError("Probe did not report successful completion")
case["passed"] = True
except RuntimeError as error:
failures += 1
case.update(passed=False, error=str(error))
cases.append(case)
(output / "results.json").write_text(json.dumps(report, indent=2) + "\n")
print(f"{'PASS' if case['passed'] else 'FAIL'} {name}", flush=True)
report["passed"] = failures == 0
report["failure_count"] = failures
(output / "results.json").write_text(json.dumps(report, indent=2) + "\n")
print(f"Installed artifact matrix: {len(cases) - failures}/{len(cases)} passed; {output / 'results.json'}")
print(f"Current 3.x artifact matrix: {len(cases) - failures}/{len(cases)} passed; {output / 'results.json'}")
return 1 if failures else 0


Expand All @@ -175,6 +167,7 @@ def main() -> int:
parser.add_argument("--output", type=Path, default=Path("target/otel-api-compatibility"))
parser.add_argument("--new-core", type=Path)
parser.add_argument("--new-plugin", type=Path)
parser.add_argument("--new-testing", type=Path)
parser.add_argument("--maven", default=shutil.which("mvn") or "mvn")
parser.add_argument("--java", default=shutil.which("java") or "java")
parser.add_argument("--javac", default=shutil.which("javac") or "javac")
Expand Down
2 changes: 1 addition & 1 deletion .github/workflows/conformance-tests.yml
Original file line number Diff line number Diff line change
Expand Up @@ -21,7 +21,7 @@ on:

concurrency:
# Runs share fixed per-suite stack names, so serialize the whole workflow
# (queue, don't cancel) to avoid concurrent CloudFormation updates on the
# with multiple pending runs to avoid concurrent CloudFormation updates on the
# same stack -- mirrors e2e-tests.yml.
group: conformance-tests
cancel-in-progress: false
Expand Down
21 changes: 21 additions & 0 deletions .github/workflows/e2e-tests.yml
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,7 @@ on:
- 'pom.xml'

concurrency:
# Shared stacks require serial deployment; preserve waiting runs from other PRs.
group: e2e-tests
cancel-in-progress: false
queue: max
Expand Down Expand Up @@ -149,6 +150,26 @@ jobs:
if [[ "$invariant_violated" == true ]]; then
exit 1
fi
- name: Collect Lambda errors after failed E2E tests
if: failure() && env.E2E_LOG_START_TIME_MS != ''
env:
E2E_STACK_NAME: Java${{ matrix.java }}-JavaSDKCloudBasedIntegrationTestStack
run: |
set -euo pipefail
log_groups=$(aws cloudformation list-stack-resources \
--stack-name "$E2E_STACK_NAME" \
--query "StackResourceSummaries[?ResourceType=='AWS::Logs::LogGroup'].PhysicalResourceId" \
--output text)
for log_group in $log_groups; do
echo "::group::Lambda errors: $log_group"
aws logs filter-log-events \
--log-group-name "$log_group" \
--start-time "$E2E_LOG_START_TIME_MS" \
--filter-pattern '%ERROR|Error|Exception|timed.out|Invalid.suspension|not.active|already.registered%' \
--limit 100 --no-paginate \
--query 'events[].message' --output text
echo "::endgroup::"
done
- name: Publish test case summary
if: always()
env:
Expand Down
4 changes: 2 additions & 2 deletions .github/workflows/otel-conformance-tests.yml
Original file line number Diff line number Diff line change
Expand Up @@ -68,14 +68,14 @@ 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@f5855f2d0f60be996973173cf479c3567f83e30f
with:
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
sdk_repository: aws/aws-durable-execution-sdk-java
sdk_ref: ${{ github.event.pull_request.head.sha || github.sha }}
conformance_test_ref: ${{ inputs.conformance_test_ref || '02d6dca971a38c13d94d6233d12f687e55b2a572' }}
conformance_test_ref: ${{ inputs.conformance_test_ref || '98b802cdb172614f98e217f7464784d47b9bb484' }}
checkout_sdk: true
# Build the handlers from this repo's checked-out module instead of the conformance repo's
# bundled examples/java. Path is relative to the conformance workspace where the SDK is
Expand Down
1 change: 1 addition & 0 deletions AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -335,6 +335,7 @@ Run `mvn spotless:apply` after Java changes. Then run the narrowest relevant tes
- [Error Handling](docs/advanced/error-handling.md)
- [Logging](docs/advanced/logging.md)
- [Migration from 1.x to 2.x](docs/migration-1.x-to-2.x.md)
- [Migration from 2.x to 3.x](docs/migration-2.x-to-3.x.md)

### Official AWS SDKs

Expand Down
1 change: 1 addition & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -119,6 +119,7 @@ See [Deploy Lambda durable functions with Infrastructure as Code](https://docs.a
- [<u>Error Handling</u>](docs/advanced/error-handling.md) - SDK exceptions for handling failures
- [<u>Logging</u>](docs/advanced/logging.md) - How to use DurableLogger
- [<u>Migrating from 1.x to 2.x</u>](docs/migration-1.x-to-2.x.md) - Upgrade guide for breaking changes since `v1.2.1`
- [<u>Migrating from 2.x to 3.x</u>](docs/migration-2.x-to-3.x.md) - Upgrade guide for the factory-only, per-invocation plugin contract
- [<u>Release Process</u>](RELEASE.md) - Prepare and publish Maven releases
- [<u>Testing</u>](docs/advanced/testing.md) - Utilities for local development and cloud-based integration testing

Expand Down
53 changes: 53 additions & 0 deletions conformance-tests-otel/README.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,53 @@
# Java OpenTelemetry conformance handlers

The shared conformance repository owns the requirements and validators. This
module supplies public-API handlers and SAM resources for both tracing views.
Existing cases 1–20 and their resources are unchanged by the additions below.

| Case | Scenario | Behavior |
| --- | --- | --- |
| 21 | `completed-step-replay` | Complete a step, suspend on a one-second durable wait, then complete another step. The first step body is skipped on replay. |
| 22 | `user-function-context` | Probe handler entry/restoration/resume, step, child, concurrent parallel branches and map iterations, and their nested steps. |
| 23 | `callback-function-context` | Probe step retry attempts, condition checks, callback submission, a wrapped retry helper's body and strategy, and a virtual child after the asynchronous work finishes. |
| 24 | `invocation-retry-status` | Throw the public retryable execution exception after a checkpointed step, then recover on replay. |

Cases 22 and 23 require a valid active `SpanContext` and create/end ordinary
`conformance.<label>` user spans with `conformance.callback=<label>`. They do not
set parents, attach contexts, copy the execution ARN onto user spans, or create
SDK spans. The shared validator determines the canonical trace from SDK spans
and checks the actual user-span parents.

Root handler probes may inherit the view's SDK root or a valid same-trace,
non-SDK ambient Lambda parent. Java currently preserves an ambient context
propagated by its configured Java agent rather than binding a new SDK root on
the handler thread. Tests must not replace that existing behavior merely to
make parent names identical across SDKs. Without ambient propagation, the
validity check exposes missing handler context. Nested callbacks must retain
their named context or attempt parents. Execution-view context placeholders
can be valid without recording, so probes do not require `isRecording()`.

Case 23 uses the public step attempt number and checkpointed condition state;
it has no mutable warm-container markers. Its helper and virtual child run
after callback completion, so normal suspension/replay does not repeat those
probes. Only the intentional helper failure is caught. The wrapped helper's
strategy executes inside its established context lifecycle; this does not
establish a parent contract for every other policy callback.

Serialization, map item naming, completion policies, standalone step retry or
condition wait strategies, and callbacks invoked before tracing setup are not
claimed to share the stable user-function hook boundary. Plugin hooks,
extractors, and samplers cannot be required to have made their future SDK span
current. Arbitrary user-created threads are also outside these cases.

Case 24 captures `isReplaying()` at handler entry before consuming the completed
step. Java's public `UnrecoverableDurableExecutionException(ErrorObject, true)`
requests invocation retry; the handler does not inject plugin hooks or
checkpoint messages. Its expected invocation statuses are `RETRYING` then
`SUCCEEDED`. Recovery may repeat telemetry delivery, so this case imposes no
exact operation-export count.

The module builds against the SDK selected by `JAVA_SDK_VERSION` or the
`durable.sdk.version` Maven property. Run matched shared requirements when
validating new cases; a workflow pinned to an older requirement set does not
validate cases 21–24. The existing configured OTel conformance workflow owns
deployment and backend validation.
Original file line number Diff line number Diff line change
@@ -0,0 +1,20 @@
// Copyright Amazon.com, Inc. or its affiliates. All Rights Reserved.
// SPDX-License-Identifier: Apache-2.0
package software.amazon.lambda.durable.conformance.otel;

import java.time.Duration;
import java.util.Map;
import software.amazon.lambda.durable.DurableContext;

/** Completed-step replay scenario for OTel requirement 21 in both views. */
public final class Otel21CompletedStepReplay extends OtelConformanceHandler<String> {

@Override
public String handleRequest(Map<String, Object> event, DurableContext context) {
requireScenario(event, "completed-step-replay");
var before = context.step("otel-before-wait", String.class, step -> "before");
context.wait("otel-replay-wait", Duration.ofSeconds(1));
var after = context.step("otel-after-wait", String.class, step -> "after");
return before + "-" + after;
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,96 @@
// Copyright Amazon.com, Inc. or its affiliates. All Rights Reserved.
// SPDX-License-Identifier: Apache-2.0
package software.amazon.lambda.durable.conformance.otel;

import io.opentelemetry.api.GlobalOpenTelemetry;
import io.opentelemetry.api.trace.Span;
import java.time.Duration;
import java.util.List;
import java.util.Map;
import software.amazon.lambda.durable.DurableContext;
import software.amazon.lambda.durable.config.MapConfig;
import software.amazon.lambda.durable.config.ParallelConfig;

/** Active user-function context scenario for OTel requirement 22 in both views. */
public final class Otel22UserFunctionContext extends OtelConformanceHandler<String> {

@Override
public String handleRequest(Map<String, Object> event, DurableContext context) {
requireScenario(event, "user-function-context");
probe("handler");
context.step("otel-context-step", String.class, step -> {
probe("step");
return "step";
});
runChild(context);
runParallel(context);
runMap(context);
probe("handler-restored");
context.wait("otel-context-resume", Duration.ofSeconds(1));
probe("handler-after-resume");
return "context-complete";
}

private static void runChild(DurableContext context) {
context.runInChildContext("otel-context-child", String.class, child -> {
probe("child");
var result = child.step("otel-context-child-step", String.class, step -> {
probe("child-step");
return "child";
});
probe("child-restored");
return result;
});
}

private static void runMap(DurableContext context) {
context.map(
"otel-context-map",
List.of(0, 1),
Integer.class,
(item, index, iteration) -> {
probe("map-" + index);
return iteration.step("otel-context-map-step-" + index, Integer.class, step -> {
probe("map-step-" + index);
return item;
});
},
MapConfig.builder()
.maxConcurrency(2)
.itemNamer((item, index) -> "otel-context-iteration-" + index)
.build());
}

private static void runParallel(DurableContext context) {
var parallel = context.parallel(
"otel-context-parallel",
ParallelConfig.builder().maxConcurrency(2).build());
try (parallel) {
parallel.branch("otel-context-branch-a", String.class, branch -> {
probe("parallel-a");
return branch.step("otel-context-branch-step-a", String.class, step -> {
probe("parallel-step-a");
return "a";
});
});
parallel.branch("otel-context-branch-b", String.class, branch -> {
probe("parallel-b");
return branch.step("otel-context-branch-step-b", String.class, step -> {
probe("parallel-step-b");
return "b";
});
});
}
}

private static void probe(String label) {
if (!Span.current().getSpanContext().isValid()) {
throw new IllegalStateException("No active span context for " + label);
}
var span = GlobalOpenTelemetry.getTracer("durable-conformance")
.spanBuilder("conformance." + label)
.setAttribute("conformance.callback", label)
.startSpan();
span.end();
}
}
Loading
Loading