Skip to content
Open
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
9 changes: 9 additions & 0 deletions packages/aws-durable-execution-sdk-python-otel/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -297,6 +297,15 @@ 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.

### Status attributes

`durable.invocation.status` uses `RETRYING` when the core plugin hook reports
`InvocationStatus.RETRY`. The core enum remains unchanged. Operation spans use
OTel `OK` only for `SUCCEEDED`, `ERROR` when error details are delivered, and
`UNSET` for other outcomes without error details, including `FAILED`,
`CANCELLED`, `TIMED_OUT`, and `STOPPED`. The original durable operation status
remains in `durable.operation.status`.

### 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 @@ -43,6 +43,7 @@
InvocationStartInfo,
OperationEndInfo,
OperationStartInfo,
OperationStatus,
OperationType,
UserFunctionEndInfo,
UserFunctionOutcome,
Expand Down Expand Up @@ -607,7 +608,9 @@ def on_invocation_end(self, info: InvocationEndInfo) -> None:
if self._invocation_span is not None:
self._invocation_span.set_attribute(
"durable.invocation.status",
info.status.value if info.status else "",
"RETRYING"
if info.status is InvocationStatus.RETRY
else info.status.value,
)
if info.status in (InvocationStatus.SUCCEEDED, InvocationStatus.PENDING):
self._invocation_span.set_status(StatusCode.OK)
Expand Down Expand Up @@ -726,7 +729,7 @@ def on_operation_end(self, info: OperationEndInfo) -> None:
span.record_exception(
Exception(info.error.message or info.error.type or "Unknown error")
)
else:
elif info.status is OperationStatus.SUCCEEDED:
span.set_status(StatusCode.OK)

self._note_parent_end(info.parent_id, end_time)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,7 @@
InvocationStartInfo,
OperationEndInfo,
OperationStartInfo,
OperationStatus,
OperationType,
UserFunctionEndInfo,
UserFunctionOutcome,
Expand Down Expand Up @@ -676,7 +677,10 @@ def on_invocation_end(self, info: InvocationEndInfo) -> None:
invocation_span = self._get_span(None)
if invocation_span:
invocation_span.set_attribute(
"durable.invocation.status", info.status.value
"durable.invocation.status",
"RETRYING"
if info.status is InvocationStatus.RETRY
else info.status.value,
)
# Span status mapping: SUCCEEDED/PENDING -> OK, FAILED -> ERROR,
# RETRY -> UNSET. RETRY is left UNSET because the plugin interface
Expand Down Expand Up @@ -780,7 +784,7 @@ def on_operation_end(self, info: OperationEndInfo) -> None:
span.record_exception(
Exception(info.error.message or info.error.type or "Unknown error")
)
else:
elif info.status is OperationStatus.SUCCEEDED:
span.set_status(StatusCode.OK)

self._end_span(info.operation_id, info.end_time)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -1123,7 +1123,9 @@ def test_invocation_span_status_kind_and_attributes(status, expected_code):
invocation = {s.name: s for s in exporter.get_finished_spans()}["Invocation"]
assert invocation.kind is trace.SpanKind.INTERNAL
assert invocation.attributes is not None
assert invocation.attributes["durable.invocation.status"] == status.value
assert invocation.attributes["durable.invocation.status"] == (
"RETRYING" if status is InvocationStatus.RETRY else status.value
)
assert invocation.attributes["durable.invocation.first"] is True
assert invocation.status.status_code is expected_code

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -563,7 +563,11 @@ def test_invocation_span_status_reflects_execution_status(
invocation = next(s for s in spans if s.name == "Invocation")
attributes = invocation.attributes
assert attributes is not None
assert attributes["durable.invocation.status"] == invocation_status.value
assert attributes["durable.invocation.status"] == (
"RETRYING"
if invocation_status is InvocationStatus.RETRY
else invocation_status.value
)
assert invocation.status.status_code is expected_span_status


Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,172 @@
"""Shared status regressions through both public plugin lifecycles."""

from collections.abc import Iterator
from datetime import UTC, datetime, timedelta
from typing import Any

import pytest
from aws_durable_execution_sdk_python.lambda_service import (
ErrorObject,
OperationStatus,
OperationSubType,
)
from aws_durable_execution_sdk_python.plugin import (
InvocationEndInfo,
InvocationStartInfo,
InvocationStatus,
OperationEndInfo,
OperationStartInfo,
OperationType,
)
from opentelemetry import context
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 opentelemetry.trace import StatusCode

from aws_durable_execution_sdk_python_otel import (
ExecutionOtelPlugin,
InvocationOtelPlugin,
OtelPluginConfig,
)


START = datetime(2026, 10, 1, tzinfo=UTC)
END = START + timedelta(seconds=1)
ARN = (
"arn:aws:lambda:us-west-2:123456789012:function:status:1/durable-execution/test/id"
)


InstrumentedView = tuple[
ExecutionOtelPlugin | InvocationOtelPlugin, InMemorySpanExporter
]


@pytest.fixture(params=[ExecutionOtelPlugin, InvocationOtelPlugin])
def instrumented_view(request: pytest.FixtureRequest) -> Iterator[InstrumentedView]:
exporter = InMemorySpanExporter()
provider = TracerProvider()
provider.add_span_processor(SimpleSpanProcessor(exporter))
plugin = request.param(
OtelPluginConfig(
tracer_provider=provider,
context_extractor=lambda _: None,
enrich_logger=False,
)
)
original = context.get_current()
yield plugin, exporter
assert context.get_current() == original
provider.shutdown()


@pytest.mark.parametrize(
"status",
[
OperationStatus.SUCCEEDED,
OperationStatus.FAILED,
OperationStatus.CANCELLED,
OperationStatus.TIMED_OUT,
OperationStatus.STOPPED,
],
)
@pytest.mark.parametrize("with_error", [False, True])
@pytest.mark.parametrize("started_here", [False, True])
def test_operation_status_and_exception(
instrumented_view: InstrumentedView,
status: OperationStatus,
with_error: bool,
started_here: bool,
) -> None:
"""A delivered terminal hook must not label an unsuccessful operation OK."""
plugin, exporter = instrumented_view
invocation: dict[str, Any] = dict(
request_id="first",
execution_arn=ARN,
execution_start_time=START,
is_first_invocation=True,
)
plugin.on_invocation_start(InvocationStartInfo(**invocation))
operation: dict[str, Any] = dict(
operation_id="wait",
name="wait",
parent_id=None,
operation_type=OperationType.WAIT,
sub_type=OperationSubType.WAIT,
start_time=START,
is_replayed=False,
)
if started_here:
plugin.on_operation_start(
OperationStartInfo(**operation, status=OperationStatus.STARTED)
)
error = (
ErrorObject(
message="deadline exceeded",
type="TimeoutError",
data=None,
stack_trace=None,
)
if with_error
else None
)
plugin.on_operation_end(
OperationEndInfo(**operation, end_time=END, status=status, error=error)
)
plugin.on_invocation_end(
InvocationEndInfo(**invocation, status=InvocationStatus.SUCCEEDED)
)
span = next(s for s in exporter.get_finished_spans() if s.name == "wait")
assert span.attributes is not None
assert span.attributes["durable.operation.status"] == status.value
expected = (
StatusCode.ERROR
if with_error
else (
StatusCode.OK if status is OperationStatus.SUCCEEDED else StatusCode.UNSET
)
)
assert span.status.status_code is expected
if with_error:
assert span.status.description == "deadline exceeded"
assert len(span.events) == 1
assert span.events[0].name == "exception"
assert span.events[0].attributes is not None
assert span.events[0].attributes["exception.message"] == "deadline exceeded"
else:
assert not span.events


def test_retry_then_resume_exports_retrying(
instrumented_view: InstrumentedView,
) -> None:
"""Telemetry normalizes RETRY without changing the core enum or next invocation."""
plugin, exporter = instrumented_view
original = context.get_current()
for index, status in enumerate(
[InvocationStatus.RETRY, InvocationStatus.SUCCEEDED]
):
invocation: dict[str, Any] = dict(
request_id=f"request-{index}",
execution_arn=ARN,
execution_start_time=START,
is_first_invocation=index == 0,
)
plugin.on_invocation_start(InvocationStartInfo(**invocation))
plugin.on_invocation_end(InvocationEndInfo(**invocation, status=status))
assert context.get_current() == original
spans = [s for s in exporter.get_finished_spans() if s.name == "Invocation"]
assert [
s.attributes["durable.invocation.status"]
for s in spans
if s.attributes is not None
] == [
"RETRYING",
"SUCCEEDED",
]
assert [s.status.status_code for s in spans] == [StatusCode.UNSET, StatusCode.OK]
assert InvocationStatus.RETRY.value == "RETRY"
assert spans[0].context is not None
assert spans[1].context is not None
assert spans[0].context.trace_id == spans[1].context.trace_id
Loading