pierrejeambrun commented on code in PR #71459:
URL: https://github.com/apache/airflow/pull/71459#discussion_r3870659439
##########
airflow-core/src/airflow/models/dagrun.py:
##########
@@ -532,10 +532,9 @@ def dag_run_data(self) -> DRDataModel:
@property
def dag_versions(self) -> list[DagVersion]:
"""Return the DAG versions associated with the TIs of this DagRun."""
- # when the dag is in a versioned bundle, we keep the dag version fixed
- if self.bundle_version:
- return [self.created_dag_version] if self.created_dag_version is
not None else []
-
+ # Always derive from actual TI/TIH versions. Do not special-case
bundle-versioned
+ # runs via created_dag_version_id: clear(...,
run_on_latest_version=True) can bump
+ # that pointer while leaving uncleared TIs on older versions
(apache/airflow#71454).
Review Comment:
```suggestion
# Always derive from actual TI/TIH versions.
```
##########
airflow-core/tests/unit/models/test_dagrun.py:
##########
@@ -1372,16 +1373,78 @@ def
test_dag_run_dag_versions_with_null_created_dag_version(self, dag_maker, ses
EmptyOperator(task_id="empty")
dag_run = dag_maker.create_dagrun()
+ ti_version_ids = {ti.dag_version_id for ti in dag_run.task_instances
if ti.dag_version_id is not None}
+ assert ti_version_ids
+
dag_run.bundle_version = "some_bundle_version"
dag_run.created_dag_version_id = None
dag_run.created_dag_version = None
session.merge(dag_run)
session.flush()
- # This should return empty list, not [None]
- assert dag_run.dag_versions == []
- assert isinstance(dag_run.dag_versions, list)
- assert len(dag_run.dag_versions) == 0
+ # Derive from TI versions; never return [None] from a null
created_dag_version shortcut.
+ versions = dag_run.dag_versions
+ assert isinstance(versions, list)
+ assert {dv.id for dv in versions} == ti_version_ids
+ assert None not in versions
+
+ def test_dag_run_dag_versions_after_partial_run_on_latest_version(self,
dag_maker, session):
+ """dag_versions must include uncleared TI versions after a partial
run_on_latest_version clear.
+
+ Repro for apache/airflow#71454: clearing a subset with
run_on_latest_version=True bumps
+ created_dag_version_id / bundle_version to latest while uncleared TIs
stay on the old
+ version. The property must report both.
Review Comment:
```suggestion
```
##########
airflow-core/tests/unit/models/test_dagrun.py:
##########
@@ -1372,16 +1373,78 @@ def
test_dag_run_dag_versions_with_null_created_dag_version(self, dag_maker, ses
EmptyOperator(task_id="empty")
dag_run = dag_maker.create_dagrun()
+ ti_version_ids = {ti.dag_version_id for ti in dag_run.task_instances
if ti.dag_version_id is not None}
+ assert ti_version_ids
+
dag_run.bundle_version = "some_bundle_version"
dag_run.created_dag_version_id = None
dag_run.created_dag_version = None
session.merge(dag_run)
session.flush()
- # This should return empty list, not [None]
- assert dag_run.dag_versions == []
- assert isinstance(dag_run.dag_versions, list)
- assert len(dag_run.dag_versions) == 0
+ # Derive from TI versions; never return [None] from a null
created_dag_version shortcut.
+ versions = dag_run.dag_versions
+ assert isinstance(versions, list)
+ assert {dv.id for dv in versions} == ti_version_ids
+ assert None not in versions
+
+ def test_dag_run_dag_versions_after_partial_run_on_latest_version(self,
dag_maker, session):
+ """dag_versions must include uncleared TI versions after a partial
run_on_latest_version clear.
+
+ Repro for apache/airflow#71454: clearing a subset with
run_on_latest_version=True bumps
+ created_dag_version_id / bundle_version to latest while uncleared TIs
stay on the old
+ version. The property must report both.
Review Comment:
We do not specify issues number in tests.
##########
airflow-core/tests/unit/models/test_dagrun.py:
##########
@@ -1372,16 +1373,78 @@ def
test_dag_run_dag_versions_with_null_created_dag_version(self, dag_maker, ses
EmptyOperator(task_id="empty")
dag_run = dag_maker.create_dagrun()
+ ti_version_ids = {ti.dag_version_id for ti in dag_run.task_instances
if ti.dag_version_id is not None}
+ assert ti_version_ids
+
dag_run.bundle_version = "some_bundle_version"
dag_run.created_dag_version_id = None
dag_run.created_dag_version = None
session.merge(dag_run)
session.flush()
- # This should return empty list, not [None]
- assert dag_run.dag_versions == []
- assert isinstance(dag_run.dag_versions, list)
- assert len(dag_run.dag_versions) == 0
+ # Derive from TI versions; never return [None] from a null
created_dag_version shortcut.
+ versions = dag_run.dag_versions
+ assert isinstance(versions, list)
+ assert {dv.id for dv in versions} == ti_version_ids
+ assert None not in versions
+
+ def test_dag_run_dag_versions_after_partial_run_on_latest_version(self,
dag_maker, session):
Review Comment:
Metion bundled version in the test name and docstring I guess
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]