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'")

Reply via email to