Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
28 commits
Select commit Hold shift + click to select a range
7d89a4e
fix: preserve handler worker context
Oct 2, 2026
0ea7ea2
test: classify handler dispatch unit coverage
Oct 2, 2026
95322ac
fix: bind missing handler trace context
zhongkechen Oct 3, 2026
812f112
fix: confine handler context to worker scopes
Oct 3, 2026
370295f
test: retain published-core context expectations
Oct 3, 2026
e1f2cc9
test: isolate the published-core OTel environment
Oct 3, 2026
bdb6aca
test: pin legacy OTel compatibility coverage
Oct 3, 2026
687f8f0
ci: verify minimum-core OTel compatibility
zhongkechen Oct 3, 2026
d459e8e
test: cover OTel handler context scopes directly
Oct 3, 2026
7d7a1a3
fix: isolate invocation plugin context bindings
zhongkechen Oct 4, 2026
d62289f
fix: preserve uninstrumented worker context
zhongkechen Oct 4, 2026
df318b7
Merge branch 'main' into fix/otel-handler-context-428
zhongkechen Oct 5, 2026
7f558a5
fix: isolate failed plugin context setup
zhongkechen Oct 6, 2026
b5cdeac
fix: require explicit handler scope opt-in
zhongkechen Oct 6, 2026
40d9edf
fix: bind execution view in the handler worker
zhongkechen Oct 6, 2026
bd11d69
Merge commit 'refs/maintenance-python-20261006/base' into maintenance…
Oct 6, 2026
788011d
Merge commit 'refs/maintenance-python-20261006/base' into maintenance…
Oct 6, 2026
66f86d3
ci: route OTel conformance to CodeBuild
Oct 7, 2026
a028769
ci: preserve queued conformance runs
Oct 7, 2026
daae092
ci: adopt shared backend queue preservation
Oct 7, 2026
74c5486
fix: build new js example workspace dependencies
Oct 7, 2026
a421382
chore: merge main after telemetry view guard
Oct 7, 2026
176026b
ci: pin published conformance workflow
Oct 7, 2026
b989dad
docs: clarify failed plugin setup context
Oct 8, 2026
d8de5f3
refactor: simplify optional handler context scopes
Oct 8, 2026
f225d88
fix(otel): unify invocation lifecycle on handler worker
Oct 8, 2026
29847f3
fix: propagate plugin context into concurrent branches
Oct 9, 2026
2456325
test: cover instrumented branch context handoff
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
101 changes: 101 additions & 0 deletions .github/scripts/install_otel_test_wheels.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,101 @@
"""Install and verify built artifacts without silently using editable sources."""

from __future__ import annotations

import argparse
import ast
import hashlib
import importlib.metadata
import json
import subprocess
import sys
import zipfile
from pathlib import Path


ROOT = Path(__file__).resolve().parents[2]
CORE = "aws-durable-execution-sdk-python"
OTEL = CORE + "-otel"


def built_wheel(package: str) -> Path:
directory = ROOT / "packages" / package
module = package.replace("-", "_")
about = ast.parse((directory / "src" / module / "__about__.py").read_text())
version = next(
ast.literal_eval(node.value)
for node in about.body
if isinstance(node, ast.Assign)
and any(
isinstance(target, ast.Name) and target.id == "__version__"
for target in node.targets
)
)
wheels = list((directory / "dist").glob(f"{module}-{version}-*.whl"))
if len(wheels) != 1:
raise ValueError(f"Build exactly one {package} {version} wheel first: {wheels}")
return wheels[0]


def verify(wheel: Path, package: str) -> None:
installed = importlib.metadata.distribution(package)
direct = json.loads(installed.read_text("direct_url.json") or "{}")
assert not direct.get("dir_info", {}).get("editable"), direct
digest = hashlib.sha256(wheel.read_bytes()).hexdigest()
assert direct["archive_info"]["hashes"]["sha256"] == digest, direct
module = package.replace("-", "_")
with zipfile.ZipFile(wheel) as archive:
sources = [
name
for name in archive.namelist()
if name.startswith(module + "/") and name.endswith(".py")
]
assert sources
for name in sources:
path = Path(installed.locate_file(name)).resolve()
assert "site-packages" in path.parts, path
assert path.read_bytes() == archive.read(name), path
print(
json.dumps(
{
"package": package,
"version": installed.version,
"wheel": str(wheel),
"sha256": digest,
"verified_sources": len(sources),
}
)
)


def main() -> None:
parser = argparse.ArgumentParser(description=__doc__)
parser.add_argument("--legacy-plugin", action="store_true")
args = parser.parse_args()
packages = [CORE] if args.legacy_plugin else [CORE, OTEL]
wheels = [built_wheel(package) for package in packages]
subprocess.run(
[
sys.executable,
"-m",
"pip",
"install",
"--no-index",
"--no-deps",
"--force-reinstall",
*map(str, wheels),
],
check=True,
)
for wheel, package in zip(wheels, packages, strict=True):
verify(wheel, package)
if args.legacy_plugin:
assert importlib.metadata.version(OTEL) == "1.0.0"
import aws_durable_execution_sdk_python_otel as otel

assert "site-packages" in Path(otel.__file__).resolve().parts
subprocess.run([sys.executable, "-m", "pip", "check"], check=True)


if __name__ == "__main__":
main()
162 changes: 162 additions & 0 deletions .github/tests/otel_lifecycle_compatibility_test.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,162 @@
"""Exercise real installed version pairs through the public durable handler."""

from __future__ import annotations

import contextvars
from datetime import UTC, datetime
from importlib.metadata import version
import os
from pathlib import Path
import threading
from types import SimpleNamespace

import pytest
from aws_durable_execution_sdk_python import durable_execution
from aws_durable_execution_sdk_python import execution as core_execution
from aws_durable_execution_sdk_python.lambda_service import (
ExecutionDetails,
Operation,
OperationStatus,
OperationType,
)
from aws_durable_execution_sdk_python.plugin import DurableInstrumentationPlugin
from aws_durable_execution_sdk_python_otel.execution_plugin import ExecutionOtelPlugin
from aws_durable_execution_sdk_python_otel.invocation_plugin import InvocationOtelPlugin
from aws_durable_execution_sdk_python_otel.otel_plugin_config import OtelPluginConfig
from opentelemetry import baggage, context, trace
from opentelemetry.sdk.trace import TracerProvider
from opentelemetry.sdk.trace.export import SimpleSpanProcessor
from opentelemetry.sdk.trace.export.in_memory_span_exporter import InMemorySpanExporter
from packaging.version import Version


class NoNetworkClient:
def __getattr__(self, name: str):
raise AssertionError(f"Unexpected service access: {name}")


@pytest.mark.parametrize("view", [InvocationOtelPlugin, ExecutionOtelPlugin])
@pytest.mark.parametrize("order", ["alone", "baggage-first", "baggage-last"])
@pytest.mark.parametrize("ambient_present", [False, True])
def test_installed_pair_preserves_context_and_documents_legacy_fallback(
view, order, ambient_present
):
legacy = os.environ.get("OTEL_COMPAT_LEGACY") == "1"
assert Version(version("aws-durable-execution-sdk-python")) >= Version("2.1.0")
assert version("aws-durable-execution-sdk-python-otel") == (
"1.0.0" if legacy else "1.1.0"
)
assert "site-packages" in Path(core_execution.__file__).resolve().parts
import aws_durable_execution_sdk_python_otel as installed_otel

assert "site-packages" in Path(installed_otel.__file__).resolve().parts

def run() -> None:
provider = TracerProvider()
exporter = InMemorySpanExporter()
provider.add_span_processor(SimpleSpanProcessor(exporter))
tracer = provider.get_tracer("installed-lifecycle")
phases = []
worker_threads = []
caller_thread = threading.get_ident()

class BaggagePlugin(DurableInstrumentationPlugin):
def on_invocation_start(self, info):
self.token = context.attach(baggage.set_baggage("customer", "present"))
phases.append("baggage-start")
worker_threads.append(threading.get_ident())

def on_invocation_end(self, info):
context.detach(self.token)
phases.append("baggage-end")
worker_threads.append(threading.get_ident())

plugin = view(
OtelPluginConfig(
tracer_provider=provider,
enrich_logger=False,
context_extractor=lambda _: None,
)
)
plugins = {
"alone": [plugin],
"baggage-first": [BaggagePlugin(), plugin],
"baggage-last": [plugin, BaggagePlugin()],
}[order]
body_parents = []

@durable_execution(plugins=plugins, boto3_client=NoNetworkClient())
def handler(event, durable_context):
phases.append("body")
worker_threads.append(threading.get_ident())
assert baggage.get_baggage("incoming") == "keep"
assert baggage.get_baggage("customer") == (
None if order == "alone" else "present"
)
parent = trace.get_current_span().get_span_context()
body_parents.append(parent)
with tracer.start_as_current_span("customer-span"):
pass
return "ok"

operation = Operation(
operation_id="installed",
operation_type=OperationType.EXECUTION,
status=OperationStatus.STARTED,
start_timestamp=datetime(2026, 10, 8, tzinfo=UTC),
execution_details=ExecutionDetails(input_payload="{}"),
)
event = {
"DurableExecutionArn": "test-arn/installed",
"CheckpointToken": "token",
"InitialExecutionState": {
"Operations": [operation.to_json_dict()],
"NextMarker": "",
},
}
lambda_context = SimpleNamespace(
aws_request_id="installed",
client_context=None,
identity=None,
_epoch_deadline_time_in_ms=0,
invoked_function_arn="test-arn",
tenant_id=None,
)
ambient = tracer.start_span("host") if ambient_present else trace.INVALID_SPAN
host = baggage.set_baggage(
"incoming", "keep", trace.set_span_in_context(ambient)
)
token = context.attach(host)
try:
for _ in range(2):
assert handler(event, lambda_context)["Status"] == "SUCCEEDED"
assert context.get_current() is host
assert caller_thread not in worker_threads
assert phases == (
["body"] * 2
if order == "alone"
else ["baggage-start", "body", "baggage-end"] * 2
)
spans = exporter.get_finished_spans()
if legacy and view is InvocationOtelPlugin:
# The released plugin does not attach an Invocation fallback.
assert body_parents == [ambient.get_span_context()] * 2
else:
name = "Invocation" if view is InvocationOtelPlugin else "Workflow"
contexts = [span.context for span in spans if span.name == name]
assert contexts
assert all(
parent.is_valid and parent in contexts for parent in body_parents
)
users = [span for span in spans if span.name == "customer-span"]
assert len(users) == 2
assert [span.parent for span in users] == [
parent if parent.is_valid else None for parent in body_parents
]
assert plugin._context_tokens == {}
finally:
context.detach(token)
ambient.end()
provider.shutdown()

contextvars.Context().run(run)
6 changes: 4 additions & 2 deletions .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -84,8 +84,6 @@ jobs:
run: hatch run types:check
- name: Run tests + coverage
run: hatch run test:cov
- name: Verify supported legacy core compatibility
run: hatch run test-pypi-otel-legacy:test
- name: Build distribution
run: |
for pkg in packages/*/; do
Expand All @@ -96,6 +94,10 @@ jobs:
cd "$GITHUB_WORKSPACE"
fi
done
- name: Test installed core and OTel wheels
run: hatch run test-wheel-otel:test
- name: Test released OTel with the new core wheel
run: hatch run test-wheel-otel-legacy:test
- name: Verify OTel wheel dependency contract
run: |
OTEL_WHEEL=$(find packages/aws-durable-execution-sdk-python-otel/dist \
Expand Down
4 changes: 3 additions & 1 deletion .github/workflows/cloud-tests.yml
Original file line number Diff line number Diff line change
Expand Up @@ -133,7 +133,9 @@ jobs:
echo "Could not resolve the latest ADOT Python layer for $AWS_REGION"
exit 1
fi
aws lambda get-layer-version-by-arn \
# Parallel jobs can throttle this read; keep retries local to the lookup.
AWS_RETRY_MODE=standard AWS_MAX_ATTEMPTS=8 \
aws lambda get-layer-version-by-arn \
--arn "$ADOT_LAYER_ARN" \
--region "$AWS_REGION" \
--query LayerVersionArn \
Expand Down
2 changes: 1 addition & 1 deletion .github/workflows/opentelemetry-conformance-tests.yml
Original file line number Diff line number Diff line change
Expand Up @@ -66,7 +66,7 @@ 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@f18bd0b5f28c5c90e288d0fb8bca08a849b51863
with:
language: python
runs_on: codebuild-github-actions-runner-${{ github.run_id }}-${{ github.run_attempt }}
Expand Down
22 changes: 16 additions & 6 deletions CONTRIBUTING.md
Original file line number Diff line number Diff line change
Expand Up @@ -70,15 +70,25 @@ hatch run dev-otel:typecheck # type check otel only
hatch run dev-examples:test # run examples tests only
```

### PyPI release testing
### Installed package compatibility testing

To verify packages work against the published PyPI version of the core SDK (rather than the local workspace):
Build the core and OTel distributions with `hatch build` in each package, then run
these commands from the repository root:

```bash
hatch run test-pypi-otel:test # test new OTel capabilities against capable installed core
hatch run test-pypi-otel-legacy:test # valid registrations/lifecycles on supported core 2.0.x
hatch run test-pypi-examples:test # test examples against PyPI core SDK
```
hatch run test-wheel-otel:test # full OTel suite on the two built wheels
hatch run test-wheel-otel-legacy:test # released OTel 1.0.0 with the built core
hatch run test-pypi-examples:test # examples against the published core
```

The wheel lanes have no editable workspace members. They verify installed source
bytes and artifact hashes before exercising public handlers. OTel 1.1 requires
the redesigned core 2.1.0 lifecycle; it no longer claims compatibility with core
2.0.x. Publish core first. Building both wheels lets CI verify the intended pair
before that minimum is available on PyPI. The legacy-plugin lane documents the
actual core-only upgrade: host isolation is provided by the new core, while old
Invocation OTel does not gain the new fallback. Workspace tests continue to cover
the complete current implementation with `hatch run dev-otel:test`.

### Package-level commands

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@
"""

import json
from threading import Lock
from typing import Any

from aws_durable_execution_sdk_python.config import Duration, ParallelConfig
Expand All @@ -31,13 +32,19 @@
)


_log_lock = Lock()


def _emit(record: dict[str, Any], execution_arn: str | None) -> None:
# Prefix every plugin record with the execution ARN as a top-level field so
# the conformance runner's CloudWatch JSON filter can scope logs to a single
# execution. Omit the field when the ARN is unset (never invent a value).
if execution_arn:
record = {"durableExecutionArn": execution_arn, **record}
print(json.dumps(record), flush=True)
# Start and end hooks can run on different threads. Keep print's separate
# body/newline writes together so the runner receives one JSON per line.
with _log_lock:
print(json.dumps(record), flush=True)


class WaitReplayFlagPlugin(DurableInstrumentationPlugin):
Expand Down
Loading
Loading