Repository navigation
fix(otel): export sampled fallback roots early #754
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
base: main
Are you sure you want to change the base?
Changes from all commits
aae44c4
3048aff
325a0c9
ff4d6ab
0c74059
9e4265c
d7eb42c
35d02c6
d657284
ffb1251
1a57d1d
dc429fd
15caff4
2492e71
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,75 @@ | ||
| """Materialize the stable SDK-owned ancestor of a fallback execution trace.""" | ||
|
|
||
| from __future__ import annotations | ||
|
|
||
| from dataclasses import dataclass | ||
| from datetime import datetime | ||
|
|
||
| from opentelemetry.context import Context | ||
| from opentelemetry.sdk.trace.sampling import SamplingResult | ||
| from opentelemetry.trace import SpanContext, SpanKind, Tracer | ||
|
|
||
| from aws_durable_execution_sdk_python_otel.deterministic_id_generator import ( | ||
| DeterministicIdGenerator, | ||
| ) | ||
| from aws_durable_execution_sdk_python_otel.durable_sampling import ( | ||
| DurableSamplingIntent, | ||
| store_sampling_intent, | ||
| ) | ||
|
|
||
|
|
||
| @dataclass(frozen=True) | ||
| class ExecutionRoot: | ||
| """A zero-duration anchor, reproducible in any invocation of an execution. | ||
|
|
||
| Re-exporting the same anchor permits recovery after a failed flush without | ||
| checkpointing telemetry state. Resources and sampler metadata remain owned | ||
| by the configured OpenTelemetry provider; SDK-added fields remain stable. | ||
| """ | ||
|
|
||
| execution_arn: str | ||
| ancestor: SpanContext | ||
| start_time: datetime | ||
|
|
||
| def export( | ||
| self, | ||
| tracer: Tracer, | ||
| id_generator: DeterministicIdGenerator, | ||
| sampling_intent: DurableSamplingIntent, | ||
| ) -> None: | ||
| """Export a sampled local ancestor; never replace an external parent.""" | ||
| if self.ancestor.is_remote or not self.ancestor.trace_flags.sampled: | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Codex AI review · Finding P1: The standalone extra still supports OpenTelemetry 1.20, whose |
||
| return | ||
|
|
||
| # Reuse the invocation's resolved decision and sampler metadata without | ||
| # resampling. An empty parent context makes this an actual root, while | ||
| # the regular tracer preserves configured resources and processors. | ||
| root_attributes: dict[str, str | bool] = { | ||
| "durable.execution.arn": self.execution_arn, | ||
| "durable.execution.synthetic_root": True, | ||
| } | ||
| # Sampler attributes normally override span attributes. Give this root | ||
| # its own intent so SDK-owned identity reaches processors intact, while | ||
| # retaining all other metadata and the invocation's original result. | ||
| result = sampling_intent.result | ||
| root_intent = DurableSamplingIntent( | ||
| SamplingResult( | ||
| result.decision, | ||
| attributes={**dict(result.attributes or {}), **root_attributes}, | ||
| trace_state=result.trace_state, | ||
| ) | ||
| ) | ||
| root_context = store_sampling_intent(Context(), root_intent) | ||
| timestamp = int(self.start_time.timestamp() * 1_000_000_000) | ||
| with id_generator._use_ids_for_span( | ||
| trace_id=self.ancestor.trace_id, | ||
| span_id=self.ancestor.span_id, | ||
| ): | ||
| span = tracer.start_span( | ||
| "DurableExecutionRoot", | ||
| context=root_context, | ||
| kind=SpanKind.INTERNAL, | ||
| attributes=root_attributes, | ||
| start_time=timestamp, | ||
| ) | ||
| span.end(end_time=timestamp) | ||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,114 @@ | ||
| """Fallback anchor coverage through the decorator and local durable runner.""" | ||
|
|
||
| import inspect | ||
| from typing import Any | ||
|
|
||
| import pytest | ||
| from aws_durable_execution_sdk_python import DurableContext | ||
| from aws_durable_execution_sdk_python.config import Duration | ||
| from aws_durable_execution_sdk_python.execution import durable_execution | ||
| from aws_durable_execution_sdk_python.plugin import ( | ||
| DurableInstrumentationPlugin, | ||
| InvocationEndInfo, | ||
| InvocationStatus, | ||
| ) | ||
| from aws_durable_execution_sdk_python_testing.runner import DurableFunctionTestRunner | ||
| from opentelemetry import context | ||
| from opentelemetry.sdk.trace import ReadableSpan, TracerProvider | ||
| from opentelemetry.sdk.trace.export import BatchSpanProcessor | ||
| from opentelemetry.sdk.trace.export.in_memory_span_exporter import InMemorySpanExporter | ||
|
|
||
| from aws_durable_execution_sdk_python_otel import ( | ||
| ExecutionOtelPlugin, | ||
| InvocationOtelPlugin, | ||
| OtelPluginConfig, | ||
| ) | ||
|
|
||
|
|
||
| class _CaptureInvocationEnd(DurableInstrumentationPlugin): | ||
| def __init__(self, exporter: InMemorySpanExporter) -> None: | ||
| self.exporter = exporter | ||
| self.snapshots: list[tuple[InvocationStatus, tuple[ReadableSpan, ...]]] = [] | ||
|
|
||
| def on_invocation_end(self, info: InvocationEndInfo) -> None: | ||
| self.snapshots.append((info.status, tuple(self.exporter.get_finished_spans()))) | ||
|
|
||
|
|
||
| @pytest.mark.parametrize("plugin_type", [ExecutionOtelPlugin, InvocationOtelPlugin]) | ||
| @pytest.mark.parametrize("outcome", ["success", "failure", "timeout"]) | ||
| def test_fallback_anchor_survives_suspension( | ||
| plugin_type: type[ExecutionOtelPlugin] | type[InvocationOtelPlugin], | ||
| outcome: str, | ||
| ) -> None: | ||
| exporter = InMemorySpanExporter() | ||
| provider = TracerProvider() | ||
| provider.add_span_processor( | ||
| BatchSpanProcessor(exporter, schedule_delay_millis=60000) | ||
| ) | ||
| plugin = plugin_type( | ||
| OtelPluginConfig( | ||
| tracer_provider=provider, | ||
| context_extractor=lambda _: None, | ||
| enrich_logger=False, | ||
| ) | ||
| ) | ||
| observer = _CaptureInvocationEnd(exporter) | ||
| completed_steps: list[str] = [] | ||
| before = context.get_current() | ||
|
|
||
| def step(_step_context: Any) -> str: | ||
| completed_steps.append("step") | ||
| return "stored" | ||
|
|
||
| def handler(_event: Any, durable: DurableContext) -> str: | ||
| value = durable.step(step, name="completed-step") | ||
| durable.wait( | ||
| Duration.from_seconds(60 if outcome == "timeout" else 1), name="wait" | ||
| ) | ||
| if outcome == "failure": | ||
| raise ValueError("failed after resume") | ||
| return value | ||
|
|
||
| wrapped = durable_execution(handler, plugins=[plugin, observer]) | ||
| try: | ||
| # Published runners use real time; newer workspace runners default to | ||
| # skipping durable waits. Exercise timeout without advancing that wait. | ||
| runner_options: dict[str, Any] = {"handler": wrapped} | ||
| if "skip_time" in inspect.signature(DurableFunctionTestRunner).parameters: | ||
| runner_options["skip_time"] = outcome != "timeout" | ||
| with DurableFunctionTestRunner(**runner_options) as runner: | ||
| arn = runner.run_async( | ||
| input="{}", timeout=3 if outcome == "timeout" else 15 | ||
| ) | ||
| result = runner.wait_for_result(arn, timeout=10) | ||
| assert completed_steps == ["step"] | ||
| assert observer.snapshots[0][0] is InvocationStatus.PENDING | ||
| first_spans = observer.snapshots[0][1] | ||
| first_roots = [s for s in first_spans if s.name == "DurableExecutionRoot"] | ||
| assert len(first_roots) == 1 | ||
| assert not any(s.name == "Workflow" for s in first_spans) | ||
| roots = [ | ||
| s for s in exporter.get_finished_spans() if s.name == "DurableExecutionRoot" | ||
| ] | ||
| assert all(s.to_json() == first_roots[0].to_json() for s in roots) | ||
| if outcome == "timeout": | ||
| assert result.status.value == "FAILED" | ||
| assert result.error is not None | ||
| assert "timed out" in (result.error.message or "") | ||
| assert len(observer.snapshots) == 1 | ||
| assert len(roots) == 1 | ||
| assert not any(s.name == "Workflow" for s in exporter.get_finished_spans()) | ||
| else: | ||
| assert result.status.value == ( | ||
| "FAILED" if outcome == "failure" else "SUCCEEDED" | ||
| ) | ||
| workflows = [ | ||
| s for s in exporter.get_finished_spans() if s.name == "Workflow" | ||
| ] | ||
| assert len(workflows) == 1 | ||
| assert workflows[0].parent is not None | ||
| assert first_roots[0].context is not None | ||
| assert workflows[0].parent.span_id == first_roots[0].context.span_id | ||
| assert context.get_current() == before | ||
| finally: | ||
| provider.shutdown() |
This comment was marked as outdated.
Sorry, something went wrong.
Uh oh!
There was an error while loading. Please reload this page.
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
Checked the published 1.20.0 release:
TraceFlags.sampledis already defined asreturn bool(self & TraceFlags.SAMPLED). The official API 1.20.0 wheel contains the same property, and this code constructs the local ancestor's flags asTraceFlagsinstances.I also installed both
opentelemetry-api==1.20.0andopentelemetry-sdk==1.20.0and ran the complete OTel suite on15caff4: all 449 tests pass, including fallback-root export in both views and decorated-handler suspend/resume cases. The existing accessor works at the declared dependency floor.