Skip to content
Open
12 changes: 12 additions & 0 deletions .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -60,6 +60,18 @@ jobs:
run: hatch run types:check
- name: Run tests + coverage
run: hatch run test:cov
- name: Test OTel with the minimum supported core
run: |
hatch run test-pypi-otel: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")
PYTHON
hatch run test-pypi-otel:test
- name: Build distribution
run: |
for pkg in packages/*/; do
Expand Down
6 changes: 5 additions & 1 deletion CONTRIBUTING.md
Original file line number Diff line number Diff line change
Expand Up @@ -75,10 +75,14 @@ hatch run dev-examples:test # run examples tests only
To verify packages work against the published PyPI version of the core SDK (rather than the local workspace):

```bash
hatch run test-pypi-otel:test # test otel against PyPI core SDK
hatch run test-pypi-otel:test # test otel against the minimum supported core (2.0.0)
hatch run test-pypi-examples:test # test examples against PyPI core SDK
```

The OTel PyPI environment excludes the local core and pins the minimum supported
release, so newer PyPI releases cannot remove legacy compatibility coverage. 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
23 changes: 23 additions & 0 deletions packages/aws-durable-execution-sdk-python-otel/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -156,6 +156,18 @@ 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.
Existing plugin registration, factory lifetime, and checkpoint formats are unchanged.

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

```python
Expand Down Expand Up @@ -297,6 +309,17 @@ The resolved decision is applied to Workflow, Invocation, operation, and attempt
spans. This avoids independently querying stateful or ratio-based samplers for
each durable span in the same invocation.

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.


### 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 @@ -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

from aws_durable_execution_sdk_python.plugin import (
Expand Down Expand Up @@ -286,12 +288,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 @@ -579,6 +577,24 @@ def on_invocation_start(self, info: InvocationStartInfo) -> None:
attributes=self._extract_attributes(info),
)

@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 @@ -655,6 +671,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
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@

from __future__ import annotations

from contextlib import nullcontext
from dataclasses import replace
from datetime import UTC, datetime
from typing import Any
Expand All @@ -25,12 +26,16 @@
OperationType,
StepDetails,
)
from aws_durable_execution_sdk_python.plugin import DurableInstrumentationPlugin
from aws_durable_execution_sdk_python_otel.deterministic_id_generator import (
derive_workflow_span_id,
)
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 context as otel_context
from opentelemetry import trace
from opentelemetry.propagators.aws.aws_xray_propagator import AwsXRayPropagator
from opentelemetry.sdk.trace import TracerProvider
from opentelemetry.sdk.trace.export import SimpleSpanProcessor
from opentelemetry.sdk.trace.export.in_memory_span_exporter import InMemorySpanExporter
Expand Down Expand Up @@ -253,3 +258,209 @@ def handler_impl(_event: Any, context: DurableContext) -> str:
assert completed_wait_span.end_time is not None
assert after_resume.start_time is not None
assert completed_wait_span.end_time <= after_resume.start_time


@pytest.mark.parametrize(
("plugin_type", "extra_context_plugin"),
[
(InvocationOtelPlugin, False),
(InvocationOtelPlugin, True),
(ExecutionOtelPlugin, False),
]
# Execution-view caller isolation requires the coordinated newer core.
# Released 2.0.x retains the pre-existing same-order teardown limitation;
# the legacy lane continues checking its supported combinations above.
+ (
[(ExecutionOtelPlugin, True)]
if hasattr(DurableInstrumentationPlugin, "handler_context")
else []
),
)
@pytest.mark.parametrize("reverse_plugins", [False, True])
@pytest.mark.parametrize("fail_after_resume", [False, True])
@pytest.mark.parametrize("ambient_kind", ["same", "unrelated", "absent"])
def test_handler_user_spans_inherit_context_across_resume_and_failure(
monkeypatch: pytest.MonkeyPatch,
plugin_type: type[InvocationOtelPlugin] | type[ExecutionOtelPlugin],
fail_after_resume: bool,
ambient_kind: str,
extra_context_plugin: bool,
reverse_plugins: bool,
) -> None:
monkeypatch.delenv("DURABLE_EXECUTION_PLUGINS", raising=False)
monkeypatch.setenv("_X_AMZN_TRACE_ID", XRAY_TRACE_HEADER)
# The documented PyPI compatibility environment deliberately uses an older
# core. Keep exercising its supported operation tracing and lifecycle while
# asserting the new handler contract only when that core exposes the scope.
supports_handler_context = hasattr(DurableInstrumentationPlugin, "handler_context")
exporter = InMemorySpanExporter()
provider = TracerProvider()
provider.add_span_processor(SimpleSpanProcessor(exporter))
plugin = plugin_type(
OtelPluginConfig(tracer_provider=provider, enrich_logger=False)
)
tracer = provider.get_tracer("customer")
before_context = otel_context.get_current()
calls: list[str] = []

def user_span(name: str) -> None:
# Ordinary instrumentation: the SDK/caller supplies the active parent.
span = tracer.start_span(name)
span.end()

def step_body(_step_context: Any) -> str:
calls.append("step")
user_span("step-user")
return "saved"

def handler_body(_event: Any, context: DurableContext) -> str:
if extra_context_plugin:
assert baggage.get_baggage("customer") == (
"present" if supports_handler_context else None
)
user_span("handler-entry")
saved = context.step(step_body, name="before-wait")
user_span("handler-after-step")
context.wait(Duration.from_seconds(1), name="context-wait")
user_span("handler-after-resume")
if fail_after_resume:
raise ValueError("handler failed after resume")
return saved

# An unrelated plugin may own a caller-thread OTel baggage scope. The
# invocation-view fallback must never become part of its saved token.
from opentelemetry import baggage

class BaggagePlugin(DurableInstrumentationPlugin):
token: Any = None

def on_invocation_start(self, _info: Any) -> None:
self.token = otel_context.attach(baggage.set_baggage("customer", "present"))

def on_invocation_end(self, _info: Any) -> None:
otel_context.detach(self.token)
self.token = None

plugins: list[DurableInstrumentationPlugin] = [plugin]
if extra_context_plugin:
plugins.append(BaggagePlugin())
if reverse_plugins:
plugins.reverse()
handler = durable_execution(handler_body, plugins=plugins)
remote = AwsXRayPropagator().extract({"X-Amzn-Trace-Id": XRAY_TRACE_HEADER})
assert trace.get_current_span(remote).get_span_context().trace_id == XRAY_TRACE_ID
initial_operations = [_execution_operation()]
checkpoint, operations = _checkpoint_store(initial_operations)
ambient_ids: list[int] = []
host_context = remote if ambient_kind == "same" else otel_context.Context()
try:
with patch(
"aws_durable_execution_sdk_python.execution.LambdaClient"
) as client_class:
client = Mock()
client.checkpoint = checkpoint
client_class.initialize_client.return_value = client
# Standard host instrumentation supplies a same-trace Lambda span.
host_scope = (
tracer.start_as_current_span("lambda-first", context=host_context)
if ambient_kind != "absent"
else nullcontext()
)
with host_scope:
host = trace.get_current_span()
ambient_ids.append(host.get_span_context().span_id)
first = handler(_event(initial_operations), _lambda_context())
assert (
trace.get_current_span().get_span_context()
== host.get_span_context()
)
assert first["Status"] == InvocationStatus.PENDING.value
assert otel_context.get_current() == before_context
resumed_operations = [
replace(
operation,
status=OperationStatus.SUCCEEDED,
end_timestamp=datetime.now(UTC),
)
if operation.name == "context-wait"
else operation
for operation in operations.values()
]
wait_id = next(
operation.operation_id
for operation in resumed_operations
if operation.name == "context-wait"
)
checkpoint, _ = _checkpoint_store(resumed_operations)
with patch(
"aws_durable_execution_sdk_python.execution.LambdaClient"
) as client_class:
client = Mock()
client.checkpoint = checkpoint
client_class.initialize_client.return_value = client
host_scope = (
tracer.start_as_current_span("lambda-resume", context=host_context)
if ambient_kind != "absent"
else nullcontext()
)
with host_scope:
host = trace.get_current_span()
ambient_ids.append(host.get_span_context().span_id)
resumed = handler(
_event(resumed_operations, updated_operation_ids=[wait_id]),
_lambda_context(),
)
assert (
trace.get_current_span().get_span_context()
== host.get_span_context()
)
assert resumed["Status"] == (
InvocationStatus.FAILED.value
if fail_after_resume
else InvocationStatus.SUCCEEDED.value
)
assert calls == ["step"]
assert otel_context.get_current() == before_context
spans = exporter.get_finished_spans()
expected_parents = (
[None, None]
if not supports_handler_context
else [derive_workflow_span_id(EXECUTION_ARN)] * 2
if plugin_type is ExecutionOtelPlugin
else ambient_ids
if ambient_kind == "same"
else [
span.context.span_id
for span in spans
if span.name == "Invocation" and span.context is not None
]
)
for name in ("handler-entry", "handler-after-step"):
users = [span for span in spans if span.name == name]
assert len(users) == 2
assert [span.parent.span_id if span.parent else None for span in users] == (
expected_parents
)
assert all(
span.context is not None
and (
(span.context.trace_id == XRAY_TRACE_ID) == supports_handler_context
)
for span in users
)
after_resume = next(
span for span in spans if span.name == "handler-after-resume"
)
assert (
after_resume.parent.span_id if after_resume.parent else None
) == expected_parents[1]
step_user = next(span for span in spans if span.name == "step-user")
assert step_user.parent is not None
assert any(
span.name == "before-wait attempt 1"
and span.context is not None
and span.context.span_id == step_user.parent.span_id
for span in spans
)
finally:
provider.shutdown()
Loading
Loading