This is an automated email from the ASF dual-hosted git repository.
kaxil pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/airflow.git
The following commit(s) were added to refs/heads/main by this push:
new e364ee76483 Keep agent framework spans when worker tracing is off
(#73903)
e364ee76483 is described below
commit e364ee76483a8f84d4e898ebc9d360856c790d72
Author: Kaxil Naik <[email protected]>
AuthorDate: Wed Sep 30 18:55:05 2026 +0100
Keep agent framework spans when worker tracing is off (#73903)
With no tracer provider installed in the worker, the task span made the
Dag run's propagated trace context current even when that context was
unsampled. A span the task's own code then started under a tracer
provider it installed, such as an agent framework's, inherited the
unsampled flag and was dropped by the default parent-based sampler.
The check is whether a provider is installed, not the [traces] otel_on
setting, so a provider installed by auto-instrumentation keeps the Dag
run's sampling decision.
---
.../src/airflow/sdk/execution_time/task_runner.py | 12 ++++
.../task_sdk/execution_time/test_task_runner.py | 64 ++++++++++++++++++++++
2 files changed, 76 insertions(+)
diff --git a/task-sdk/src/airflow/sdk/execution_time/task_runner.py
b/task-sdk/src/airflow/sdk/execution_time/task_runner.py
index 5dc3de05e13..ed622023308 100644
--- a/task-sdk/src/airflow/sdk/execution_time/task_runner.py
+++ b/task-sdk/src/airflow/sdk/execution_time/task_runner.py
@@ -177,6 +177,18 @@ def _make_task_span(msg: StartupDetails):
parent_context = (
TraceContextTextMapPropagator().extract(msg.ti.context_carrier) if
msg.ti.context_carrier else None
)
+ if (
+ parent_context is not None
+ and not
trace.get_current_span(parent_context).get_span_context().trace_flags.sampled
+ and isinstance(trace.get_tracer_provider(),
(trace.ProxyTracerProvider, trace.NoOpTracerProvider))
+ ):
+ # With no tracer provider installed, the no-op tracer of
opentelemetry-api 1.40 and later
+ # still makes the propagated context current. When that context is
unsampled, every span
+ # the task's own code starts under a tracer provider it installs (an
agent framework's,
+ # for example) inherits the "not sampled" flag and a parent-based
sampler, the
+ # OpenTelemetry default, drops it. Leave such spans as roots. A
provider installed before
+ # the task started, by core tracing or by auto-instrumentation,
samples as it was set up.
+ parent_context = None
ti = msg.ti
span_name = f"worker.{ti.task_id}"
if ti.map_index is not None and ti.map_index >= 0:
diff --git a/task-sdk/tests/task_sdk/execution_time/test_task_runner.py
b/task-sdk/tests/task_sdk/execution_time/test_task_runner.py
index ce9b1187c5e..9cdb75446b1 100644
--- a/task-sdk/tests/task_sdk/execution_time/test_task_runner.py
+++ b/task-sdk/tests/task_sdk/execution_time/test_task_runner.py
@@ -788,6 +788,70 @@ def
test_task_span_no_parent_when_no_context_carrier(make_ti_context):
assert finished[0].parent is None
+# The trace and parent ids from the W3C Trace Context specification's
traceparent example; the trailing
+# flags byte is 00, "not sampled".
+UNSAMPLED_TRACE_ID = 0x4BF92F3577B34DA6A3CE929D0E0E4736
+UNSAMPLED_CARRIER = {"traceparent":
f"00-{UNSAMPLED_TRACE_ID:032x}-00f067aa0ba902b7-00"}
+
+
+def _make_startup_with_carrier(make_ti_context, carrier: dict[str, str]) ->
StartupDetails:
+ return StartupDetails(
+ ti=TaskInstance(
+ id=uuid7(),
+ task_id="my_task",
+ dag_id="test_dag",
+ run_id="test_run",
+ try_number=1,
+ dag_version_id=uuid7(),
+ queue="default",
+ context_carrier=carrier,
+ ),
+ dag_rel_path="",
+ bundle_info=BundleInfo(name="my-bundle", version=None),
+ ti_context=make_ti_context(),
+ start_date=timezone.utcnow(),
+ sentry_integration="",
+ )
+
+
+def
test_with_no_tracer_provider_an_unsampled_parent_is_not_made_current(make_ti_context):
+ """A span the task code starts under its own provider must not inherit
"not sampled"."""
+ task_code_exporter = InMemorySpanExporter()
+ task_code_provider = TracerProvider()
+
task_code_provider.add_span_processor(SimpleSpanProcessor(task_code_exporter))
+
+ with (
+ mock.patch("airflow.sdk.execution_time.task_runner.tracer",
trace.NoOpTracer()),
+ mock.patch.object(trace, "get_tracer_provider",
return_value=trace.ProxyTracerProvider()),
+ _make_task_span(_make_startup_with_carrier(make_ti_context,
UNSAMPLED_CARRIER)),
+ ):
+ current = trace.get_current_span().get_span_context()
+ with
task_code_provider.get_tracer("agent_framework").start_as_current_span("agent
run"):
+ pass
+
+ assert not current.is_valid
+ recorded = task_code_exporter.get_finished_spans()
+ assert [span.name for span in recorded] == ["agent run"]
+ assert recorded[0].parent is None
+
+
+def
test_with_a_tracer_provider_installed_the_dag_runs_sampling_decision_holds(make_ti_context):
+ """Core tracing or auto-instrumentation installed a provider: the run was
sampled out, so stay out."""
+ installed = TracerProvider()
+ task_code_exporter = InMemorySpanExporter()
+ installed.add_span_processor(SimpleSpanProcessor(task_code_exporter))
+
+ with (
+ mock.patch("airflow.sdk.execution_time.task_runner.tracer",
installed.get_tracer("airflow")),
+ mock.patch.object(trace, "get_tracer_provider",
return_value=installed),
+ _make_task_span(_make_startup_with_carrier(make_ti_context,
UNSAMPLED_CARRIER)),
+ ):
+ with
installed.get_tracer("agent_framework").start_as_current_span("agent run"):
+ pass
+
+ assert task_code_exporter.get_finished_spans() == ()
+
+
def test_parse_module_in_bundle_root(tmp_path: Path, make_ti_context):
"""Check that the bundle path is added to sys.path, so Dags can import
shared modules."""
tmp_path.joinpath("util.py").write_text("NAME = 'dag_name'")