This is an automated email from the ASF dual-hosted git repository.

shahar1 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 854a9d8e6a3 Align standard operator test modules with the source 
layout (#73229)
854a9d8e6a3 is described below

commit 854a9d8e6a3857cdaeddee66dd83ae47cf8352da
Author: Keith <[email protected]>
AuthorDate: Wed Sep 16 14:17:10 2026 +0900

    Align standard operator test modules with the source layout (#73229)
    
    The guard list reports test_branch.py and test_latest_only.py as
    missing, but both modules are tested — the files just carry an
    _operator suffix the source layout does not mirror, so the guard cannot
    see them. Rename them to the mirroring paths, drop the two stale
    exemptions, and close the two coverage gaps the audit surfaced:
    choose_branch declaring NotImplementedError and a branch operator
    returning None skipping every direct downstream task.
---
 .../tests/unit/always/test_project_structure.py    |  2 -
 .../{test_branch_operator.py => test_branch.py}    | 63 ++++++++++++++++++++++
 ...latest_only_operator.py => test_latest_only.py} |  0
 3 files changed, 63 insertions(+), 2 deletions(-)

diff --git a/airflow-core/tests/unit/always/test_project_structure.py 
b/airflow-core/tests/unit/always/test_project_structure.py
index 2b6610e2e10..19c8620a984 100644
--- a/airflow-core/tests/unit/always/test_project_structure.py
+++ b/airflow-core/tests/unit/always/test_project_structure.py
@@ -153,9 +153,7 @@ class TestProjectStructure:
             "providers/google/tests/unit/google/test_go_module_utils.py",
             
"providers/microsoft/azure/tests/unit/microsoft/azure/operators/test_adls.py",
             
"providers/snowflake/tests/unit/snowflake/triggers/test_snowflake_trigger.py",
-            "providers/standard/tests/unit/standard/operators/test_branch.py",
             "providers/standard/tests/unit/standard/operators/test_empty.py",
-            
"providers/standard/tests/unit/standard/operators/test_latest_only.py",
             
"providers/standard/tests/unit/standard/sensors/test_external_task.py",
             "providers/sftp/tests/unit/sftp/test_exceptions.py",
         ]
diff --git 
a/providers/standard/tests/unit/standard/operators/test_branch_operator.py 
b/providers/standard/tests/unit/standard/operators/test_branch.py
similarity index 85%
rename from 
providers/standard/tests/unit/standard/operators/test_branch_operator.py
rename to providers/standard/tests/unit/standard/operators/test_branch.py
index e05d3863bce..9048e31ecc4 100644
--- a/providers/standard/tests/unit/standard/operators/test_branch_operator.py
+++ b/providers/standard/tests/unit/standard/operators/test_branch.py
@@ -64,6 +64,16 @@ class ChooseBranchThree(BaseBranchOperator):
         return ["branch_3"]
 
 
+class ChooseNoneBranch(BaseBranchOperator):
+    def choose_branch(self, context):
+        return None
+
+
+def test_choose_branch_is_abstract():
+    with pytest.raises(NotImplementedError):
+        BaseBranchOperator(task_id="make_choice").choose_branch({})
+
+
 class TestBranchOperator:
     def test_without_dag_run(self, dag_maker):
         """This checks the defensive against non-existent tasks in a dag run"""
@@ -204,6 +214,59 @@ class TestBranchOperator:
                 else:
                     raise Exception
 
+    def test_none_branch_skips_all_downstream(self, dag_maker):
+        dag_id = "branch_operator_test"
+        triggered_by_kwargs = {"triggered_by": DagRunTriggeredByType.TEST} if 
AIRFLOW_V_3_0_PLUS else {}
+        with dag_maker(
+            dag_id,
+            default_args={"owner": "airflow", "start_date": DEFAULT_DATE},
+            schedule=INTERVAL,
+            serialized=True,
+        ):
+            branch_1 = EmptyOperator(task_id="branch_1")
+            branch_2 = EmptyOperator(task_id="branch_2")
+            branch_op = ChooseNoneBranch(task_id="make_choice")
+            branch_1.set_upstream(branch_op)
+            branch_2.set_upstream(branch_op)
+        if AIRFLOW_V_3_0_1:
+            dr = dag_maker.create_dagrun(
+                run_type=DagRunType.MANUAL,
+                start_date=timezone.utcnow(),
+                logical_date=DEFAULT_DATE,
+                state=State.RUNNING,
+                data_interval=DataInterval(DEFAULT_DATE, DEFAULT_DATE),
+                **triggered_by_kwargs,
+            )
+
+            with pytest.raises(DownstreamTasksSkipped) as exc_info:
+                dag_maker.run_ti("make_choice", dr)
+
+            assert sorted(exc_info.value.tasks) == [("branch_1", -1), 
("branch_2", -1)]
+        else:
+            dr = dag_maker.create_dagrun(
+                run_type=DagRunType.MANUAL,
+                start_date=timezone.utcnow(),
+                execution_date=DEFAULT_DATE,
+                state=State.RUNNING,
+                data_interval=DataInterval(DEFAULT_DATE, DEFAULT_DATE),
+                **triggered_by_kwargs,
+            )
+
+            dag_maker.run_ti("make_choice", dr)
+
+            expected = {
+                "make_choice": State.SUCCESS,
+                "branch_1": State.SKIPPED,
+                "branch_2": State.SKIPPED,
+            }
+
+            ti_date = TI.logical_date if AIRFLOW_V_3_0_PLUS else 
TI.execution_date
+
+            for ti in dag_maker.session.scalars(
+                select(TI).where(TI.dag_id == dag_id, ti_date == DEFAULT_DATE)
+            ):
+                assert ti.state == expected[ti.task_id]
+
     def test_with_skip_in_branch_downstream_dependencies(self, dag_maker):
         dag_id = "branch_operator_test"
         triggered_by_kwargs = {"triggered_by": DagRunTriggeredByType.TEST} if 
AIRFLOW_V_3_0_PLUS else {}
diff --git 
a/providers/standard/tests/unit/standard/operators/test_latest_only_operator.py 
b/providers/standard/tests/unit/standard/operators/test_latest_only.py
similarity index 100%
rename from 
providers/standard/tests/unit/standard/operators/test_latest_only_operator.py
rename to providers/standard/tests/unit/standard/operators/test_latest_only.py

Reply via email to