Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
22 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
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
13 changes: 13 additions & 0 deletions .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -86,6 +86,19 @@ jobs:
run: hatch run test:cov
- name: Verify supported legacy core compatibility
run: hatch run test-pypi-otel-legacy:test
- name: Test OTel with the minimum supported core
run: |
hatch run test-pypi-otel-minimum:python - <<'PYTHON'
from importlib.metadata import version
from pathlib import Path
import aws_durable_execution_sdk_python.plugin as plugin

assert version("aws-durable-execution-sdk-python") == "2.0.0"
assert "site-packages" in Path(plugin.__file__).parts
assert not hasattr(plugin.DurableInstrumentationPlugin, "handler_context")
assert not hasattr(plugin, "DURABLE_INSTRUMENTATION_HANDLER_CONTEXT_API_VERSION")
PYTHON
hatch run test-pypi-otel-minimum:test
- name: Build distribution
run: |
for pkg in packages/*/; do
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
8 changes: 8 additions & 0 deletions CONTRIBUTING.md
Original file line number Diff line number Diff line change
Expand Up @@ -77,9 +77,17 @@ To verify packages work against the published PyPI version of the core SDK (rath
```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-otel-minimum:test # all prior OTel tests and valid registrations on core 2.0.0
hatch run test-pypi-examples:test # test examples against PyPI core SDK
```

The OTel minimum-core environment excludes the local core and pins 2.0.0, so
newer PyPI releases cannot remove legacy compatibility coverage. It retains all
pre-existing OTel tests plus valid registration/wait-resume cases. Exclusivity
validation requires the newer core and is exercised by the complete workspace
suite and the capable-core environment. Use `hatch run dev-otel:test` for the
current workspace core.

### Package-level commands

Some commands still run from within a package directory:
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
Original file line number Diff line number Diff line change
@@ -0,0 +1,72 @@
"""Concurrent plugin callbacks must emit separate parseable JSON records."""

from __future__ import annotations

import importlib.util
import json
from concurrent.futures import ThreadPoolExecutor
from pathlib import Path
from threading import Event
from typing import Any

import pytest

import aws_durable_execution_sdk_python.execution as execution


def test_concurrent_wait_hooks_keep_complete_stdout_records(
monkeypatch: pytest.MonkeyPatch,
) -> None:
# Import the real fixture without constructing a Lambda client: this test
# exercises its stdout producer, not a deployed durable invocation.
monkeypatch.setattr(execution, "durable_execution", lambda **_: lambda fn: fn)
path = (
Path(__file__).resolve().parents[1]
/ "handlers/plugin/plugin_wait_replay_flag.py"
)
spec = importlib.util.spec_from_file_location("wait_replay_log_fixture", path)
assert spec is not None and spec.loader is not None
module = importlib.util.module_from_spec(spec)
spec.loader.exec_module(module)

first_body = Event()
second_ready = Event()
second_done = Event()
chunks: list[str] = []

class FragmentingStdout:
def write(self, text: str) -> int:
chunks.append(text)
if '"operation-start"' in text:
first_body.set()
# print writes its body and newline separately. Permit the
# other real hook to run between them unless _emit serializes it.
second_done.wait(0.2)
return len(text)

def flush(self) -> None:
pass

records: list[dict[str, Any]] = [
{"plugin": "CONFPLUGIN", "hook": "operation-start", "name": "long"},
{"plugin": "CONFPLUGIN", "hook": "operation-end", "name": "short"},
]

def emit_end() -> None:
second_ready.set()
assert first_body.wait(2)
module._emit(records[1], "execution-arn")
second_done.set()

with monkeypatch.context() as capture:
capture.setattr("sys.stdout", FragmentingStdout())
with ThreadPoolExecutor(max_workers=2) as executor:
end = executor.submit(emit_end)
assert second_ready.wait(2)
start = executor.submit(module._emit, records[0], "execution-arn")
start.result(timeout=2)
end.result(timeout=2)
actual = [json.loads(line) for line in "".join(chunks).splitlines() if line]
assert actual == [
{"durableExecutionArn": "execution-arn", **record} for record in records
]
37 changes: 37 additions & 0 deletions packages/aws-durable-execution-sdk-python-otel/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -180,6 +180,27 @@ lambda_.Function(
)
```

### Handler context propagation

A core SDK with handler-worker context propagation carries the context established
by invocation-start hooks into the handler. Invocation view preserves an active
ambient span on the canonical execution trace; when that context is absent or
belongs to a different trace, its optional `handler_context` scope makes the Invocation
span current only while the handler runs. The scope closes on the same worker in
reverse plugin order, including on failure and suspension, without changing the
invocation-hook caller. Older cores ignore this optional scope and retain their
existing behavior; install the updated core as well to get handler context propagation.
The bundled OTel classes explicitly opt in with `__durable_handler_context_api__ = 1`.
Custom subclasses must repeat that literal marker on their own concrete class to
use the scope; inherited or instance markers are ignored. Unopted legacy helpers
and properties with the same name are never inspected. Older cores ignore the
marker without importing any new core API.
Execution view similarly restores the Workflow span inside the handler scope if
another invocation-start hook clears the active span or switches to an unrelated
trace. Both views retain valid same-trace parents and baggage, and restore the
worker's previous context when the scope ends.
Existing plugin registration, factory lifetime, and checkpoint formats are unchanged.

### 3. In your Lambda handler (index.py)

```python
Expand Down Expand Up @@ -330,6 +351,22 @@ OTel `OK` only for `SUCCEEDED`, `ERROR` when error details are delivered, and
`CANCELLED`, `TIMED_OUT`, and `STOPPED`. The original durable operation status
remains in `durable.operation.status`.

### Invocation context isolation

Invocation hooks retain their caller thread and registration order. With the
updated core, invocation-local context-variable bindings are isolated from the
host: hooks see the incoming context and the handler receives their resulting
context, while invocation exit restores the host's original bindings even if a
plugin fails during setup or cleanup. Plugins must not use invocation context
bindings to mutate the host context after the invocation has returned. Older
supported cores retain their existing lifecycle behavior, including the
execution-view limitation when later plugins open invocation context scopes.
The new isolation applies only when plugins are registered. If an invocation-start
hook or handler-scope entry raises, subsequent setup and the handler retain the
bindings from before that hook. Successful scopes still clean up in their original
context, preserving token ownership. This isolates context-variable bindings;
it does not undo a plugin's mutations to shared objects or external side effects.

### Log Correlation

When `enrich_logger=True` (the default), the plugin installs a logging filter on
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -31,9 +31,11 @@

from __future__ import annotations

import contextlib
import datetime
import logging
import threading
from collections.abc import Iterator
from typing import Any, ClassVar

from aws_durable_execution_sdk_python.plugin import (
Expand Down Expand Up @@ -127,6 +129,7 @@ class ExecutionOtelPlugin(DurableInstrumentationPlugin):
span).
"""

__durable_handler_context_api__ = 1
Comment thread
zhongkechen marked this conversation as resolved.
__durable_registration_api__: ClassVar[int] = 1
exclusive_group: ClassVar[str | None] = "aws-durable-execution-otel-view"

Expand Down Expand Up @@ -509,6 +512,24 @@ def on_invocation_start(self, info: InvocationStartInfo) -> None:
if self._config.enrich_logger:
install_log_filter(self)

@contextlib.contextmanager
def handler_context(self, info: InvocationStartInfo) -> Iterator[None]:
"""Keep handler instrumentation on this execution's trace in its worker."""
ambient = trace.get_current_span().get_span_context()
workflow = self._workflow_span
token = None
if (
self._tracing_enabled
and workflow is not None
and (not ambient.is_valid or ambient.trace_id != self._execution_trace_id)
):
token = otel_context.attach(trace.set_span_in_context(workflow))
try:
yield
finally:
if token is not None:
otel_context.detach(token)

def _start_workflow_span(self, info: InvocationStartInfo) -> None:
"""Install a non-recording placeholder for the execution-scoped Workflow span.

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -2,9 +2,11 @@

from __future__ import annotations

import contextlib
import datetime
import logging
import threading
from collections.abc import Iterator
from typing import Any, ClassVar

from aws_durable_execution_sdk_python.plugin import (
Expand Down Expand Up @@ -108,6 +110,8 @@ class InvocationOtelPlugin(DurableInstrumentationPlugin):
provider installed by the ADOT Lambda layer).
"""

__durable_handler_context_api__ = 1

DEFAULT_INSTRUMENT_NAME = "aws-durable-execution-sdk-python"

__durable_registration_api__: ClassVar[int] = 1
Expand Down Expand Up @@ -299,12 +303,8 @@ def get_current_span_context(self) -> SpanContext | None:
context this is the active context span (attached in
on_user_function_start). Unrelated ambient spans are ignored so logs
stay correlated to the durable execution trace.
2. The invocation span from the plugin registry. This is the path used
for top-level handler code: the invocation span is never attached to
the worker thread's context, so the registry is the only way to
resolve it. It also covers code between top-level operations, where
detaching the operation scope restores a context with no durable
span.
2. The invocation span from the plugin registry, including lifecycle
phases outside the optional handler-worker context scope.

Returns:
A valid SpanContext, or None if no span is active.
Expand Down Expand Up @@ -597,6 +597,24 @@ def on_invocation_start(self, info: InvocationStartInfo) -> None:
if self._enrich_logger:
install_log_filter(self)

@contextlib.contextmanager
def handler_context(self, info: InvocationStartInfo) -> Iterator[None]:
"""Bind the fallback only inside the SDK-owned handler worker scope."""
ambient = trace.get_current_span().get_span_context()
invocation_span = self._get_span(None)
token = None
if (
self._tracing_enabled
and invocation_span is not None
and (not ambient.is_valid or ambient.trace_id != self._execution_trace_id)
):
token = context.attach(trace.set_span_in_context(invocation_span))
try:
yield
finally:
if token is not None:
context.detach(token)
Comment thread
zhongkechen marked this conversation as resolved.

def _start_workflow_span(self, info: InvocationStartInfo) -> None:
"""Install a non-recording placeholder for the execution-scoped Workflow span.

Expand Down Expand Up @@ -673,6 +691,10 @@ def on_invocation_end(self, info: InvocationEndInfo) -> None:
self._reset_state()
return

# User execution has finished; the worker has already closed its handler
# context scope without modifying the invocation-hook caller.
self._detach_remaining_contexts()

# Spans are registered parent-first, so close pending spans in reverse
# order to keep every child contained within its parent.
with self._operation_spans_lock:
Expand Down
Loading
Loading