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
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,10 @@
ancestor: a propagated backend parent when present, otherwise a deterministic
synthetic root. The Workflow span is exported exactly once, when the execution
reaches a terminal status. Operations are parented under the Workflow span (or
their parent operation) and *linked* to the current Invocation span.
their parent operation) and *linked* to the current Invocation span. Each
operation is likewise exported exactly once, on its deterministic span ID, when
it reaches a terminal status; while it spans invocations it is held as a
non-recording placeholder so no recording span is abandoned.

This is the Python adaptation of the JS ``ExecutionOtelPlugin`` from
aws-durable-execution-sdk-js#729. Because the Python plugin interface differs
Expand Down Expand Up @@ -56,6 +59,7 @@
SpanKind,
StatusCode,
Tracer,
TraceState,
)

from aws_durable_execution_sdk_python_otel.context_extractors import (
Expand Down Expand Up @@ -133,12 +137,16 @@ def __init__(self, config: OtelPluginConfig | None = None) -> None:
# Per-invocation state.
self._execution_arn = ""
self._execution_trace_id: int | None = None
self._execution_start_time: datetime.datetime | None = None
self._extracted_context: ExtractedContext | None = None
self._execution_trace_context: ExecutionTraceContext | None = None
self._sampling_intent: DurableSamplingIntent | None = None
self._workflow_span: Span | None = None
self._invocation_span: Span | None = None
self._operation_spans: dict[str, Span] = {}
# Operations whose span was already exported this invocation, so a
# repeated on_operation_end does not export it twice.
self._ended_operation_ids: set[str] = set()
# Tokens returned by context.attach(), keyed by the span registry key,
# paired with the thread that attached them. Every attach the plugin
# owns is released through _detach_context so the plugin never leaves a
Expand Down Expand Up @@ -292,6 +300,39 @@ def _resolve_parent(self, parent_id: str | None) -> Span | None:
return existing
return self._workflow_span

def _resolved_trace_state(self) -> TraceState:
"""Return the resolved sampling trace state, else the ancestor state.

The sampling result preserves a same-trace ambient ``tracestate`` that
the empty ancestor state would drop.
"""
intent = self._sampling_intent
if intent is not None and intent.result.trace_state is not None:
return intent.result.trace_state
if self._execution_trace_context is not None:
return self._execution_trace_context.execution_ancestor.trace_state
return TraceState()

def _operation_span_context(self, operation_id: str) -> SpanContext | None:
"""Return the deterministic SpanContext for a logical operation."""
execution_trace_context = self._execution_trace_context
if execution_trace_context is None:
return None
return SpanContext(
trace_id=execution_trace_context.trace_id,
span_id=operation_id_to_span_id(self._execution_arn, operation_id),
is_remote=False,
trace_flags=execution_trace_context.trace_flags,
trace_state=self._resolved_trace_state(),
)

def _register_operation_placeholder(self, operation_id: str) -> None:
"""Register a non-recording placeholder holding the operation context."""
span_context = self._operation_span_context(operation_id)
if span_context is None:
return
self._set_span(operation_id, NonRecordingSpan(span_context))

def _invocation_parent_context(self) -> Context:
"""Return same-trace ambient context, else execution ancestor context."""
execution_trace_context = self._execution_trace_context
Expand Down Expand Up @@ -346,6 +387,7 @@ def on_invocation_start(self, info: InvocationStartInfo) -> None:
)
self._tracing_enabled = False
return
self._execution_start_time = info.execution_start_time
self._extracted_context = _ensure_extracted_context(
self._context_extractor(info)
)
Expand Down Expand Up @@ -393,11 +435,40 @@ def on_invocation_start(self, info: InvocationStartInfo) -> None:
)

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

The Workflow span spans the whole durable execution and is exported once,
on the terminal invocation. During every invocation the plugin only needs
its deterministic SpanContext -- to parent operation spans, to keep the
Workflow current so auto-instrumented spans join the execution trace, and
for log correlation. A non-recording placeholder fills that role so a
non-terminal invocation never abandons a recording span. The recording
span is created and ended once by :meth:`_export_workflow_span`.
"""
if not self._execution_arn:
logger.warning("No execution ARN; skipping Workflow span creation")
return
if self._execution_trace_context is None:
return
workflow_span_context = SpanContext(
trace_id=self._execution_trace_context.trace_id,
span_id=derive_workflow_span_id(self._execution_arn),
is_remote=False,
trace_flags=self._execution_trace_context.trace_flags,
trace_state=self._resolved_trace_state(),
)
self._workflow_span = NonRecordingSpan(workflow_span_context)

def _export_workflow_span(self, info: InvocationEndInfo) -> None:
"""Create and end the recording Workflow span once, on a terminal status.

Uses the same deterministic span ID as the placeholder and the shared
execution ancestor as its parent, so the exported Workflow span stays on
the execution trace and correlates with every operation span across all
invocations. Anchored at the execution start time.
"""
if not self._execution_arn or self._execution_trace_context is None:
return
parent_context = self._with_sampling(
trace.set_span_in_context(
NonRecordingSpan(self._execution_trace_context.execution_ancestor),
Expand All @@ -408,13 +479,44 @@ def _start_workflow_span(self, info: InvocationStartInfo) -> None:
trace_id=None,
span_id=derive_workflow_span_id(self._execution_arn),
):
self._workflow_span = self._tracer.start_span(
workflow_span = self._tracer.start_span(
name=self._workflow_span_name,
kind=SpanKind.INTERNAL,
attributes={"durable.execution.arn": self._execution_arn},
start_time=_to_otel_timestamp(info.execution_start_time),
attributes={
"durable.execution.arn": self._execution_arn,
"durable.execution.status": (
info.status.value if info.status else ""
),
},
start_time=_to_otel_timestamp(self._execution_start_time),
context=parent_context,
)
if info.status is InvocationStatus.FAILED:
workflow_span.set_status(
StatusCode.ERROR, info.error.message if info.error else ""
)
elif info.status is InvocationStatus.SUCCEEDED:
workflow_span.set_status(StatusCode.OK)
workflow_span.end()

def _end_open_recording_spans(self) -> None:
"""End recording user-function spans left open by a suspended operation.

Operation placeholders are non-recording and export their span from
on_operation_end, so they are skipped. Reverse order keeps each child
contained within its parent; the invocation span is ended by the caller.
"""
with self._lock:
keys = list(reversed(self._operation_spans))
for key in keys:
if key == _INVOCATION_KEY:
continue
span = self._get_span(key)
if span is None or not span.is_recording():
continue
popped = self._pop_span(key)
if popped is not None:
popped.end()
Comment thread
ayushiahjolia marked this conversation as resolved.
Comment thread
ayushiahjolia marked this conversation as resolved.
Comment thread
zhongkechen marked this conversation as resolved.
Comment thread
ayushiahjolia marked this conversation as resolved.

def _start_invocation_span(self, info: InvocationStartInfo) -> None:
self._invocation_span = self._tracer.start_span(
Expand All @@ -434,12 +536,6 @@ def on_invocation_end(self, info: InvocationEndInfo) -> None:
self._reset_state()
return

# Operation spans still open here belong to operations that suspended
# (e.g. PENDING/RETRYING) rather than completed this invocation. They are
# ended only by on_operation_end; drop the references without ending them
# so they are not exported as if completed. _reset_state
# clears the span map below.

# End the invocation span regardless of terminal status. Record the
# invocation status and map it to a span status:
# SUCCEEDED/PENDING -> OK (this invocation did its work, whether it
Expand All @@ -462,23 +558,16 @@ def on_invocation_end(self, info: InvocationEndInfo) -> None:
)
self._invocation_span.end()

# The Workflow span (execution view) is exported only on a terminal
# status; otherwise its reference is dropped without ending it. Its span
# status reflects the execution outcome: SUCCEEDED -> OK, FAILED -> ERROR
# (RETRY/PENDING are non-terminal and never reach here -> UNSET).
if self._workflow_span is not None:
if info.status in _TERMINAL_INVOCATION_STATUSES:
self._workflow_span.set_attribute(
"durable.execution.status",
info.status.value if info.status else "",
)
if info.status is InvocationStatus.FAILED:
self._workflow_span.set_status(
StatusCode.ERROR, info.error.message if info.error else ""
)
elif info.status is InvocationStatus.SUCCEEDED:
self._workflow_span.set_status(StatusCode.OK)
self._workflow_span.end()
# End recording user-function spans left open by a suspended operation.
self._end_open_recording_spans()

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Codex AI review · Finding arf_v1_vsgiwjg6y5l4vgkuagwxcho636

[P1] Keep suspended CONTEXT operations non-recording. on_user_function_start() still creates a recording, deterministic CONTEXT span, so this cleanup exports it when the invocation becomes PENDING. Replay then creates and exports another span with the same trace/span ID when the context completes, allowing collectors to overwrite or reject the authoritative terminal span. Use a non-recording CONTEXT placeholder until on_operation_end() materializes the sole recording span, and add a child-context suspend/replay test asserting one terminal export.


# The Workflow span (execution view) is a non-recording placeholder
# during the invocation, so only a terminal status materializes and ends
# the recording span. Its span status reflects the execution outcome:
# SUCCEEDED -> OK, FAILED -> ERROR (RETRY/PENDING are non-terminal and
# leave the Workflow span unexported until a later terminal invocation).
if info.status in _TERMINAL_INVOCATION_STATUSES:
self._export_workflow_span(info)

self._reset_state()

Expand All @@ -495,10 +584,12 @@ def _reset_state(self) -> None:
self._extracted_context = None
self._execution_trace_context = None
self._sampling_intent = None
self._execution_start_time = None
self._workflow_span = None
self._invocation_span = None
with self._lock:
self._operation_spans = {}
self._ended_operation_ids = set()
self._tracing_enabled = False

# ------------------------------------------------------------------
Expand All @@ -510,23 +601,28 @@ def on_operation_start(self, info: OperationStartInfo) -> None:
return
if info.operation_type is OperationType.CONTEXT:
return # tracked via on_user_function_start
parent = self._resolve_parent(info.parent_id)
self._start_span(
operation_id=info.operation_id,
name=info.name or info.operation_id,
info=info,
parent=parent,
start_time=info.start_time,
)
# Hold a non-recording placeholder while the operation is open; its
# recording span is exported once on terminal on_operation_end.
self._register_operation_placeholder(info.operation_id)

def on_operation_end(self, info: OperationEndInfo) -> None:
logger.debug("Durable operation ended: %s", info)
if not self._tracing_enabled:
return
# Export the span only on the first end for this operation.
with self._lock:
if info.operation_id in self._ended_operation_ids:
return
self._ended_operation_ids.add(info.operation_id)
# Export the operation's single recording span with the deterministic ID.
span = self._get_span(info.operation_id)
if span is None:
# Cross-invocation stitching: operation started in a prior
# invocation. Create + immediately end a linked span.
if span is not None and span.is_recording():
# A CONTEXT whose user function ran this invocation already has its
# recording span; reuse it.
span.set_attributes(self._operation_attributes(info))
else:
# A placeholder, or an operation with no local span: create it now.
self._pop_span(info.operation_id)
parent = self._resolve_parent(info.parent_id)
span = self._start_span(
operation_id=info.operation_id,
Expand All @@ -535,8 +631,6 @@ def on_operation_end(self, info: OperationEndInfo) -> None:
parent=parent,
start_time=info.start_time,
)
else:
span.set_attributes(self._operation_attributes(info))

if info.error:
span.set_status(StatusCode.ERROR, info.error.message or "")
Expand Down Expand Up @@ -564,7 +658,11 @@ def _start_span(
span_key: str | None = None,
deterministic: bool = True,
) -> Span:
"""Start a span for an operation/attempt and register it."""
"""Start a recording span for an operation/attempt and register it.

Operation spans use the deterministic operation span ID; attempt spans
pass ``deterministic=False`` for a fresh ID beneath the operation span.
"""
key = span_key if span_key is not None else operation_id
with self._lock:
links = self._build_invocation_links()
Expand Down
Loading
Loading