kaxil commented on code in PR #74365:
URL: https://github.com/apache/airflow/pull/74365#discussion_r4201174907
##########
airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_tasks.py:
##########
@@ -278,6 +279,50 @@ def test_should_respond_200_serialized(self, test_client,
testing_dag_bundle):
assert response.status_code == 200
assert response.json() == expected
+ def test_dag_version_id_describes_the_task_as_it_was(self, test_client,
dag_maker, session):
+ """A later edit must not change how an earlier version's task is
reported."""
+ dag_id = "test_versioned_task_dag"
+
+ # A new DagVersion is cut only when the serialized Dag changes under a
new bundle version,
+ # so each version is written the way the shared multi-version fixture
does it.
+ with dag_maker(dag_id, session=session, bundle_version="commit-one"):
+ EmptyOperator(task_id=self.task_id, retries=4)
+ session.commit()
+ first_version_id = str(DagVersion.get_version(dag_id,
session=session).id)
+
+ with dag_maker(dag_id, session=session, bundle_version="commit-two"):
+ EmptyOperator(task_id=self.task_id, retries=0)
+ session.commit()
+
+ assert str(DagVersion.get_version(dag_id, session=session).id) !=
first_version_id
+
+ test_client.app.dependency_overrides[dag_bag_from_app] = DBDagBag
+ url = f"{self.api_prefix}/{dag_id}/tasks/{self.task_id}"
+
+ latest = test_client.get(url)
+ assert latest.status_code == 200
+ assert latest.json()["retries"] == 0
+
+ pinned = test_client.get(url, params={"dag_version_id":
first_version_id})
+ assert pinned.status_code == 200
+ assert pinned.json()["retries"] == 4
+
+ def test_dag_version_id_of_another_dag_is_not_found(self, test_client,
testing_dag_bundle):
Review Comment:
Both new tests use an id that resolves to a real serialized Dag, so the
`versioned_dag is None` half of the check isn't covered. If `versioned_dag is
None or` were dropped, an unknown id would hit `AttributeError` and return 500,
and both tests would still pass. Could you parametrize this with a random
`uuid.uuid4()` expecting 404 as well?
##########
airflow-core/src/airflow/api_fastapi/core_api/routes/public/tasks.py:
##########
@@ -102,9 +103,31 @@ def get_tasks(
),
dependencies=[Depends(requires_access_dag(method="GET",
access_entity=DagAccessEntity.TASK))],
)
-def get_task(dag_id: str, task_id, session: SessionDep, dag_bag: DagBagDep) ->
TaskResponse:
- """Get simplified representation of a task."""
- dag = get_latest_version_of_dag(dag_bag, dag_id, session)
+def get_task(
+ dag_id: str,
+ task_id,
+ session: SessionDep,
+ dag_bag: DagBagDep,
+ dag_version_id: UUID | None = None,
Review Comment:
The other public endpoints that select a version take the per-Dag
`version_number` (`get_dag_source`'s `?version_number=`,
`/dagVersions/{version_number}`, the TI `version_number` filter). Would
`version_number: int | None` fit better here? Resolving it with
`DagVersion.get_version(dag_id, version_number, session=session)` scopes the
lookup to the path Dag, so the cross-Dag check below goes away. The mapped TI
details tab, which already calls `useTaskServiceGetTask`, only has
`dag_version_number` from `LightGridTaskInstanceSummary` in its outlet context.
If the UUID is deliberate (it's the DagBag cache key and saves a lookup),
could you say so in the PR? This would be the first public version selector of
that shape, so it sets the precedent.
##########
airflow-core/src/airflow/api_fastapi/core_api/routes/public/tasks.py:
##########
@@ -102,9 +103,31 @@ def get_tasks(
),
dependencies=[Depends(requires_access_dag(method="GET",
access_entity=DagAccessEntity.TASK))],
)
-def get_task(dag_id: str, task_id, session: SessionDep, dag_bag: DagBagDep) ->
TaskResponse:
- """Get simplified representation of a task."""
- dag = get_latest_version_of_dag(dag_bag, dag_id, session)
+def get_task(
+ dag_id: str,
+ task_id,
+ session: SessionDep,
+ dag_bag: DagBagDep,
+ dag_version_id: UUID | None = None,
+) -> TaskResponse:
+ """
+ Get simplified representation of a task.
+
+ ``dag_version_id`` pins the lookup to one Dag version, so a caller
describing a past task
Review Comment:
This docstring becomes the OpenAPI description (it's already in the yaml and
the generated TS), so it's read by API consumers. "Taken per task instance
rather than per run" refers to inputs this endpoint doesn't have, and it
doesn't say what omitting it does or when it 404s. Maybe something like: "Pass
`dag_version_id` (e.g. a task instance's `dag_version.id`) to get the task as
defined in that Dag version. Without it the latest version is used. Returns 404
if the version doesn't exist or belongs to another Dag."
##########
airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_tasks.py:
##########
@@ -278,6 +279,50 @@ def test_should_respond_200_serialized(self, test_client,
testing_dag_bundle):
assert response.status_code == 200
assert response.json() == expected
+ def test_dag_version_id_describes_the_task_as_it_was(self, test_client,
dag_maker, session):
+ """A later edit must not change how an earlier version's task is
reported."""
+ dag_id = "test_versioned_task_dag"
+
+ # A new DagVersion is cut only when the serialized Dag changes under a
new bundle version,
Review Comment:
`dag_maker._make_serdag` writes a new `DagVersion` whenever the serialized
hash changes, whatever `bundle_version` is, so this comment overstates the
bundle version's role. The assert on L297 already guards the precondition, so
I'd drop the comment. The `DBDagBag` override on L299 (and L318) also repeats
what the autouse `setup` does on L74.
--
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]