Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
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
189 changes: 189 additions & 0 deletions .github/scripts/verify_otel_api_compatibility.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,189 @@
#!/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.

Dependency resolution and every probe are required: a network, compilation, or case
failure returns nonzero. No production dependency versions are modified.
"""
from __future__ import annotations

import argparse
import hashlib
import json
import os
from pathlib import Path
import shutil
import subprocess
import sys
import xml.etree.ElementTree as ET

API_VERSIONS = ("1.49.0", "1.66.0")
RELEASED_VERSION = "2.2.1"
CORE = "aws-durable-execution-sdk-java"
PLUGIN = "aws-durable-execution-sdk-java-plugin-otel"
PROBE = "software.amazon.lambda.durable.otel.InstalledApiProbe"


def execute(command: list[str], log: Path, *, env: dict[str, str] | None = None, timeout: int = 300) -> None:
with log.open("w") as output:
try:
result = subprocess.run(command, stdout=output, stderr=subprocess.STDOUT, env=env,
check=False, timeout=timeout)
except subprocess.TimeoutExpired as error:
raise RuntimeError(f"Command timed out after {timeout}s; log={log}") from error
if result.returncode:
tail = "\n".join(log.read_text(errors="replace").splitlines()[-35:])
raise RuntimeError(f"Command failed ({result.returncode}); log={log}\n{tail}")


def artifact(entries: list[Path], name: str, version: str) -> Path:
matches = [p for p in entries if p.name == f"{name}-{version}.jar"]
if len(matches) != 1:
raise RuntimeError(f"Expected exactly one {name}:{version}, got {matches}")
if not matches[0].is_file():
raise RuntimeError(f"Resolved artifact is absent: {matches[0]}")
return matches[0].resolve()


def jar_facts(path: Path) -> dict[str, str]:
return {"path": str(path), "sha256": hashlib.sha256(path.read_bytes()).hexdigest()}


def snapshot_candidate(path: Path, output: Path) -> Path:
expected = jar_facts(path)["sha256"]
directory = output / "candidate-artifacts"
directory.mkdir(exist_ok=True)
target = directory / path.name
if path.resolve() != target.resolve():
shutil.copyfile(path, target)
if jar_facts(target)["sha256"] != expected or jar_facts(path)["sha256"] != expected:
raise RuntimeError(f"Candidate changed while being snapshotted: {path}")
return target.resolve()


def candidate_jar(root: Path, module: str, name: str) -> Path:
pom = ET.parse(root / "pom.xml")
version = pom.findtext("{http://maven.apache.org/POM/4.0.0}version")
if not version:
raise RuntimeError("Cannot resolve the reactor version from pom.xml")
path = root / module / "target" / f"{name}-{version}.jar"
if not path.is_file():
raise RuntimeError(f"Build the candidate first; artifact missing: {path}")
return path.resolve()


def resolve_classpaths(fixture: Path, output: Path, maven: str) -> dict[str, list[Path]]:
classpaths: dict[str, list[Path]] = {}
for version in API_VERSIONS:
target = output / f"dependencies-{version}.txt"
execute([
maven, "-B", "-f", str(fixture / "pom.xml"),
"org.apache.maven.plugins:maven-dependency-plugin:3.11.0:build-classpath",
f"-Dotel.api.version={version}", f"-Dmdep.outputFile={target}",
], output / f"resolve-{version}.log")
classpaths[version] = [Path(p).resolve() for p in target.read_text().strip().split(os.pathsep)]
artifact(classpaths[version], "opentelemetry-api", version)
artifact(classpaths[version], "opentelemetry-context", version)
return classpaths


def probe_environment(view: str) -> dict[str, str]:
env = os.environ.copy()
# The fixture sets its own plugin registration/global provider. Do not inherit
# Lambda-hosted CI tracing or a developer's auto-agent/plugin configuration.
for key in ("_X_AMZN_TRACE_ID", "DURABLE_EXECUTION_PLUGINS", "JAVA_TOOL_OPTIONS",
"JDK_JAVA_OPTIONS", "OTEL_JAVAAGENT_EXTENSIONS", "AWS_LAMBDA_EXEC_WRAPPER"):
env.pop(key, None)
env["DURABLE_EXECUTION_PLUGINS"] = f"{view},compat-healthy"
return env


def run_matrix(args: argparse.Namespace) -> int:
root = args.root.resolve()
output = args.output.resolve()
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)
classes = output / "classes"
classes.mkdir(exist_ok=True)
execute([args.javac, "--release", "17", "-classpath", os.pathsep.join(map(str, cp["1.66.0"])),
"-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),
"candidate_inputs": candidate_inputs,
"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)
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'}")
return 1 if failures else 0


def main() -> int:
parser = argparse.ArgumentParser(description=__doc__)
parser.add_argument("--root", type=Path, default=Path(__file__).resolve().parents[2])
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("--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")
try:
return run_matrix(parser.parse_args())
except (RuntimeError, OSError, subprocess.SubprocessError) as error:
print(f"Compatibility harness failed: {error}", file=sys.stderr)
return 1


if __name__ == "__main__":
raise SystemExit(main())
5 changes: 4 additions & 1 deletion .github/workflows/ai-pr-review-address.yml
Original file line number Diff line number Diff line change
Expand Up @@ -23,7 +23,10 @@ jobs:
contents: read
issues: read
pull-requests: read
uses: aws/aws-durable-execution-ci/.github/workflows/ai-pr-review-address.yml@d6b017da14385908951d23e26c790b28a4e5f9f0
uses: aws/aws-durable-execution-ci/.github/workflows/ai-work-item-resolver.yml@d6b017da14385908951d23e26c790b28a4e5f9f0
with:
work-scope: review
upload-work-items: true

address:
if: >-
Expand Down
8 changes: 8 additions & 0 deletions .github/workflows/build.yml
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,8 @@ on:
- 'sdk-testing/**'
- 'sdk-integration-tests/**'
- 'insight-plugin/**'
- 'otel-plugin/**'
- '.github/scripts/verify_otel_api_compatibility.py'
- 'examples/**'
- 'pom.xml'
push:
Expand All @@ -42,6 +44,8 @@ on:
- 'sdk-testing/**'
- 'sdk-integration-tests/**'
- 'insight-plugin/**'
- 'otel-plugin/**'
- '.github/scripts/verify_otel_api_compatibility.py'
- 'examples/**'
- 'pom.xml'

Expand Down Expand Up @@ -81,6 +85,10 @@ jobs:
- name: Build and test
run: mvn -B install --file pom.xml

- name: Verify installed OpenTelemetry API compatibility
Comment thread
zhongkechen marked this conversation as resolved.
if: ${{ matrix.java == 17 }}
run: python3 .github/scripts/verify_otel_api_compatibility.py

- name: Setup uv for coverage badge
if: ${{ matrix.java == 17 }}
# cicirello/jacoco-badge-generator is a Docker-based action; the
Expand Down
1 change: 1 addition & 0 deletions .github/workflows/conformance-tests.yml
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,7 @@ concurrency:
# same stack -- mirrors e2e-tests.yml.
group: conformance-tests
cancel-in-progress: false
queue: max

permissions:
contents: read
Expand Down
1 change: 1 addition & 0 deletions .github/workflows/e2e-tests.yml
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,7 @@ on:
concurrency:
group: e2e-tests
cancel-in-progress: false
queue: max

# permission can be added at job level or workflow level
permissions:
Expand Down
39 changes: 35 additions & 4 deletions .github/workflows/otel-conformance-tests.yml
Original file line number Diff line number Diff line change
Expand Up @@ -53,6 +53,12 @@ on:

permissions: {}

# Serialize shared test resources and retain pending PR runs instead of replacing them.
concurrency:
group: otel-conformance-tests
cancel-in-progress: false
queue: max

jobs:
opentelemetry:
# 1. Run for non-PR events, such as scheduled runs and manual invocations
Expand All @@ -62,8 +68,9 @@ jobs:
actions: write
contents: read
id-token: write
uses: aws/aws-durable-execution-conformance-tests/.github/workflows/opentelemetry-orchestrator.yml@43872f8d5dd917fe3ee6d01a6c953c1ff67c5fe4
uses: aws/aws-durable-execution-conformance-tests/.github/workflows/opentelemetry-orchestrator.yml@a66037abbbfa55fde97f714e30f0bc262edefd63
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
Expand All @@ -75,13 +82,37 @@ jobs:
# checked out (.build/durable-sdk).
examples_dir: .build/durable-sdk/conformance-tests-otel
setup_command: |
if [ -z "${JAVA_HOME_21_X64:-}" ]; then
echo "The runner does not provide Java 21"
set -euo pipefail
if [ -n "${JAVA_HOME_21_X64:-}" ]; then
export JAVA_HOME="$JAVA_HOME_21_X64"
elif [ -x /usr/lib/jvm/java-21-amazon-corretto/bin/java ]; then
export JAVA_HOME=/usr/lib/jvm/java-21-amazon-corretto
elif [ -z "${JAVA_HOME:-}" ]; then
export JAVA_HOME="$(dirname "$(dirname "$(readlink -f "$(command -v java)")")")"
fi
if [ ! -x "$JAVA_HOME/bin/java" ] || [ ! -x "$JAVA_HOME/bin/javac" ]; then
echo "Selected JAVA_HOME=$JAVA_HOME is not a complete JDK"
exit 1
fi
export JAVA_HOME="$JAVA_HOME_21_X64"
export PATH="$JAVA_HOME/bin:$PATH"
java_spec=$(java -XshowSettings:properties -version 2>&1 | awk '$1 == "java.specification.version" {print $3}')
if [ "$java_spec" != "21" ]; then
echo "Java 21 is required; selected JAVA_HOME=$JAVA_HOME reports specification version $java_spec"
exit 1
fi
javac_version=$(javac -version 2>&1)
if [[ ! "$javac_version" =~ ^javac[[:space:]]21([.]|[[:space:]]|$) ]]; then
echo "Java 21 javac is required; selected compiler reports $javac_version"
exit 1
fi
java -version
javac -version
if [ -n "${GITHUB_ENV:-}" ]; then
echo "JAVA_HOME=$JAVA_HOME" >> "$GITHUB_ENV"
fi
if [ -n "${GITHUB_PATH:-}" ]; then
echo "$JAVA_HOME/bin" >> "$GITHUB_PATH"
fi
prepare_command: |
JAVA_SDK_VERSION=$(
mvn -B -q \
Expand Down
13 changes: 13 additions & 0 deletions otel-plugin/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -151,6 +151,19 @@ public class MyHandler extends DurableHandler<MyInput, MyOutput> {
}
```

### OpenTelemetry version compatibility

Keep the OpenTelemetry API, context, SDK, and Java agent versions aligned. This plugin is built and tested against
OpenTelemetry 1.66.0. Global-provider binding needs `GlobalOpenTelemetry.isSet()` and `getOrNoop()`; when the visible
API lacks either method (for example, API 1.49.0), the plugin logs a compatibility diagnostic and disables its telemetry
for that invocation. It does not install a no-op global that would prevent a provider from being registered later.

The existing 2.x plugin constructors, registration interfaces, and instance lifetime are retained. Nonfatal linkage
errors from plugin callbacks are logged and isolated so healthy plugins and the handler can continue. Fatal JVM errors
and `ThreadDeath` retain their existing propagation behavior. Provider registration and configuration validation remain
unchanged. Align incompatible dependencies to restore instrumentation; error isolation does not make every old
agent/API combination capable of exporting telemetry.

### 4. Grant Permissions

The function's execution role needs the `AWSXRayDaemonWriteAccess` managed policy (or equivalent permissions) to write traces to X-Ray.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -151,6 +151,19 @@ record ProviderSetup(SdkTracerProvider sdkTracerProvider, Tracer tracer) {}
* @return the resolved provider and tracer, or {@code null} when telemetry must be disabled for this invocation
*/
static ProviderSetup tryResolveGlobalProvider(String instrumentationName, String pluginName) {
try {
return resolveGlobalProvider(instrumentationName, pluginName);
} catch (LinkageError error) {
logger.warn(
"{} telemetry is disabled for this invocation because the visible OpenTelemetry dependencies "
+ "are incompatible. Align the OpenTelemetry API, SDK, and Java agent versions.",
pluginName,
error);
return null;
}
}

private static ProviderSetup resolveGlobalProvider(String instrumentationName, String pluginName) {
if (!OtelPluginAutoConfigurationState.isInstalled()) {
logger.warn(
"{} telemetry is disabled for this invocation because "
Expand All @@ -160,6 +173,9 @@ static ProviderSetup tryResolveGlobalProvider(String instrumentationName, String
javaAgentExtensionsDiagnostic());
return null;
}
if (!supportsGlobalProviderLookup(pluginName)) {
return null;
}
if (!GlobalOpenTelemetry.isSet()) {
logger.warn(
"{} telemetry is disabled for this invocation because GlobalOpenTelemetry is not initialized yet. "
Expand All @@ -186,6 +202,22 @@ static ProviderSetup tryResolveGlobalProvider(String instrumentationName, String
getSdkTracerProviderForFlush(tracerProvider, pluginName), tracerProvider.get(instrumentationName));
}

private static boolean supportsGlobalProviderLookup(String pluginName) {
try {
GlobalOpenTelemetry.class.getMethod("isSet");
GlobalOpenTelemetry.class.getMethod("getOrNoop");
return true;
} catch (NoSuchMethodException missingApi) {
logger.warn(
"{} telemetry is disabled for this invocation because the visible OpenTelemetry API lacks {}. "
+ "Global provider binding requires GlobalOpenTelemetry.isSet() and getOrNoop(); "
+ "align the API, SDK, and Java agent versions.",
pluginName,
missingApi.getMessage());
return false;
}
}

/** Returns the SdkTracerProvider for flushing, or null if the provider is wrapped by the agent classloader. */
static SdkTracerProvider getSdkTracerProviderForFlush(TracerProvider tracerProvider, String pluginName) {
if (tracerProvider instanceof SdkTracerProvider sdkTracerProvider) {
Expand Down
Loading
Loading