This is an automated email from the ASF dual-hosted git repository.
bbovenzi 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 05dec276c72 Let the task endpoint describe a task at a given Dag
version (#74365)
05dec276c72 is described below
commit 05dec276c7227b74ae8b121c7c69c4543b36100b
Author: Brent Bovenzi <[email protected]>
AuthorDate: Wed Oct 7 16:18:32 2026 -0400
Let the task endpoint describe a task at a given Dag version (#74365)
A task instance is explained by the Dag as it was when that instance ran,
but the endpoint always answered from the newest version. Anything
describing a past instance therefore reported a task that may since have
been edited, renamed or had its asset declarations changed.
The version is selected by number, matching the other public version
selectors, which also scopes the lookup to the Dag in the path. Omitting it
keeps the previous behaviour.
---
.../core_api/openapi/v2-rest-api-generated.yaml | 19 +++++++++++-
.../api_fastapi/core_api/routes/public/tasks.py | 29 ++++++++++++++++--
.../src/airflow/ui/openapi-gen/queries/common.ts | 5 ++--
.../ui/openapi-gen/queries/ensureQueryData.ts | 10 +++++--
.../src/airflow/ui/openapi-gen/queries/prefetch.ts | 10 +++++--
.../src/airflow/ui/openapi-gen/queries/queries.ts | 10 +++++--
.../src/airflow/ui/openapi-gen/queries/suspense.ts | 10 +++++--
.../ui/openapi-gen/requests/services.gen.ts | 8 +++++
.../airflow/ui/openapi-gen/requests/types.gen.ts | 1 +
.../core_api/routes/public/test_tasks.py | 35 ++++++++++++++++++++++
10 files changed, 123 insertions(+), 14 deletions(-)
diff --git
a/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml
b/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml
index ba08ddd4fe6..1f293c4cf7c 100644
---
a/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml
+++
b/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml
@@ -9886,7 +9886,16 @@ paths:
tags:
- Task
summary: Get Task
- description: Get simplified representation of a task.
+ description: 'Get simplified representation of a task.
+
+
+ Pass ``version_number`` (e.g. a task instance''s
``dag_version.version_number``)
+ to get the task
+
+ as defined in that Dag version. Without it the latest version is used.
Returns
+ 404 if that
+
+ version of the Dag does not exist.'
operationId: get_task
security:
- OAuth2PasswordBearer: []
@@ -9903,6 +9912,14 @@ paths:
required: true
schema:
title: Task Id
+ - name: version_number
+ in: query
+ required: false
+ schema:
+ anyOf:
+ - type: integer
+ - type: 'null'
+ title: Version Number
responses:
'200':
description: Successful Response
diff --git
a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/tasks.py
b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/tasks.py
index 0f95d2fc401..753b8cb0cee 100644
--- a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/tasks.py
+++ b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/tasks.py
@@ -29,6 +29,7 @@ from airflow.api_fastapi.core_api.datamodels.tasks import
TaskCollectionResponse
from airflow.api_fastapi.core_api.openapi.exceptions import
create_openapi_http_exception_doc
from airflow.api_fastapi.core_api.security import requires_access_dag
from airflow.exceptions import TaskNotFound
+from airflow.models.dag_version import DagVersion
tasks_router = AirflowRouter(tags=["Task"], prefix="/dags/{dag_id}/tasks")
@@ -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,
+ version_number: int | None = None,
+) -> TaskResponse:
+ """
+ Get simplified representation of a task.
+
+ Pass ``version_number`` (e.g. a task instance's
``dag_version.version_number``) to get the task
+ as defined in that Dag version. Without it the latest version is used.
Returns 404 if that
+ version of the Dag does not exist.
+ """
+ if version_number is None:
+ dag = get_latest_version_of_dag(dag_bag, dag_id, session)
+ else:
+ dag_version = DagVersion.get_version(dag_id, version_number,
session=session)
+ versioned_dag = None if dag_version is None else
dag_bag.get_dag(dag_version.id, session=session)
+ if versioned_dag is None:
+ raise HTTPException(
+ status.HTTP_404_NOT_FOUND,
+ f"The Dag {dag_id}, version_number {version_number} was not
found",
+ )
+ dag = versioned_dag
try:
task = dag.get_task(task_id=task_id)
except TaskNotFound:
diff --git a/airflow-core/src/airflow/ui/openapi-gen/queries/common.ts
b/airflow-core/src/airflow/ui/openapi-gen/queries/common.ts
index ae7d2972fc6..2da2837f163 100644
--- a/airflow-core/src/airflow/ui/openapi-gen/queries/common.ts
+++ b/airflow-core/src/airflow/ui/openapi-gen/queries/common.ts
@@ -958,10 +958,11 @@ export const UseTaskServiceGetTasksKeyFn = ({ dagId,
orderBy }: {
export type TaskServiceGetTaskDefaultResponse = Awaited<ReturnType<typeof
TaskService.getTask>>;
export type TaskServiceGetTaskQueryResult<TData =
TaskServiceGetTaskDefaultResponse, TError = unknown> = UseQueryResult<TData,
TError>;
export const useTaskServiceGetTaskKey = "TaskServiceGetTask";
-export const UseTaskServiceGetTaskKeyFn = ({ dagId, taskId }: {
+export const UseTaskServiceGetTaskKeyFn = ({ dagId, taskId, versionNumber }: {
dagId: string;
taskId: unknown;
-}, queryKey?: Array<unknown>) => [useTaskServiceGetTaskKey, ...(queryKey ?? [{
dagId, taskId }])];
+ versionNumber?: number;
+}, queryKey?: Array<unknown>) => [useTaskServiceGetTaskKey, ...(queryKey ?? [{
dagId, taskId, versionNumber }])];
export type VariableServiceGetVariableDefaultResponse =
Awaited<ReturnType<typeof VariableService.getVariable>>;
export type VariableServiceGetVariableQueryResult<TData =
VariableServiceGetVariableDefaultResponse, TError = unknown> =
UseQueryResult<TData, TError>;
export const useVariableServiceGetVariableKey = "VariableServiceGetVariable";
diff --git a/airflow-core/src/airflow/ui/openapi-gen/queries/ensureQueryData.ts
b/airflow-core/src/airflow/ui/openapi-gen/queries/ensureQueryData.ts
index 01b17dee670..1bda7c5964f 100644
--- a/airflow-core/src/airflow/ui/openapi-gen/queries/ensureQueryData.ts
+++ b/airflow-core/src/airflow/ui/openapi-gen/queries/ensureQueryData.ts
@@ -1893,16 +1893,22 @@ export const ensureUseTaskServiceGetTasksData =
(queryClient: QueryClient, { dag
/**
* Get Task
* Get simplified representation of a task.
+*
+* Pass ``version_number`` (e.g. a task instance's
``dag_version.version_number``) to get the task
+* as defined in that Dag version. Without it the latest version is used.
Returns 404 if that
+* version of the Dag does not exist.
* @param data The data for the request.
* @param data.dagId
* @param data.taskId
+* @param data.versionNumber
* @returns TaskResponse Successful Response
* @throws ApiError
*/
-export const ensureUseTaskServiceGetTaskData = (queryClient: QueryClient, {
dagId, taskId }: {
+export const ensureUseTaskServiceGetTaskData = (queryClient: QueryClient, {
dagId, taskId, versionNumber }: {
dagId: string;
taskId: unknown;
-}) => queryClient.ensureQueryData({ queryKey:
Common.UseTaskServiceGetTaskKeyFn({ dagId, taskId }), queryFn: () =>
TaskService.getTask({ dagId, taskId }) });
+ versionNumber?: number;
+}) => queryClient.ensureQueryData({ queryKey:
Common.UseTaskServiceGetTaskKeyFn({ dagId, taskId, versionNumber }), queryFn:
() => TaskService.getTask({ dagId, taskId, versionNumber }) });
/**
* Get Variable
* Get a variable entry.
diff --git a/airflow-core/src/airflow/ui/openapi-gen/queries/prefetch.ts
b/airflow-core/src/airflow/ui/openapi-gen/queries/prefetch.ts
index f19316f7263..58a8c63b91d 100644
--- a/airflow-core/src/airflow/ui/openapi-gen/queries/prefetch.ts
+++ b/airflow-core/src/airflow/ui/openapi-gen/queries/prefetch.ts
@@ -1893,16 +1893,22 @@ export const prefetchUseTaskServiceGetTasks =
(queryClient: QueryClient, { dagId
/**
* Get Task
* Get simplified representation of a task.
+*
+* Pass ``version_number`` (e.g. a task instance's
``dag_version.version_number``) to get the task
+* as defined in that Dag version. Without it the latest version is used.
Returns 404 if that
+* version of the Dag does not exist.
* @param data The data for the request.
* @param data.dagId
* @param data.taskId
+* @param data.versionNumber
* @returns TaskResponse Successful Response
* @throws ApiError
*/
-export const prefetchUseTaskServiceGetTask = (queryClient: QueryClient, {
dagId, taskId }: {
+export const prefetchUseTaskServiceGetTask = (queryClient: QueryClient, {
dagId, taskId, versionNumber }: {
dagId: string;
taskId: unknown;
-}) => queryClient.prefetchQuery({ queryKey:
Common.UseTaskServiceGetTaskKeyFn({ dagId, taskId }), queryFn: () =>
TaskService.getTask({ dagId, taskId }) });
+ versionNumber?: number;
+}) => queryClient.prefetchQuery({ queryKey:
Common.UseTaskServiceGetTaskKeyFn({ dagId, taskId, versionNumber }), queryFn:
() => TaskService.getTask({ dagId, taskId, versionNumber }) });
/**
* Get Variable
* Get a variable entry.
diff --git a/airflow-core/src/airflow/ui/openapi-gen/queries/queries.ts
b/airflow-core/src/airflow/ui/openapi-gen/queries/queries.ts
index 50aa7a65187..ca724d5fcb7 100644
--- a/airflow-core/src/airflow/ui/openapi-gen/queries/queries.ts
+++ b/airflow-core/src/airflow/ui/openapi-gen/queries/queries.ts
@@ -1893,16 +1893,22 @@ export const useTaskServiceGetTasks = <TData =
Common.TaskServiceGetTasksDefault
/**
* Get Task
* Get simplified representation of a task.
+*
+* Pass ``version_number`` (e.g. a task instance's
``dag_version.version_number``) to get the task
+* as defined in that Dag version. Without it the latest version is used.
Returns 404 if that
+* version of the Dag does not exist.
* @param data The data for the request.
* @param data.dagId
* @param data.taskId
+* @param data.versionNumber
* @returns TaskResponse Successful Response
* @throws ApiError
*/
-export const useTaskServiceGetTask = <TData =
Common.TaskServiceGetTaskDefaultResponse, TError = unknown, TQueryKey extends
Array<unknown> = unknown[]>({ dagId, taskId }: {
+export const useTaskServiceGetTask = <TData =
Common.TaskServiceGetTaskDefaultResponse, TError = unknown, TQueryKey extends
Array<unknown> = unknown[]>({ dagId, taskId, versionNumber }: {
dagId: string;
taskId: unknown;
-}, queryKey?: TQueryKey, options?: Omit<UseQueryOptions<TData, TError>,
"queryKey" | "queryFn">) => useQuery<TData, TError>({ queryKey:
Common.UseTaskServiceGetTaskKeyFn({ dagId, taskId }, queryKey), queryFn: () =>
TaskService.getTask({ dagId, taskId }) as TData, ...options });
+ versionNumber?: number;
+}, queryKey?: TQueryKey, options?: Omit<UseQueryOptions<TData, TError>,
"queryKey" | "queryFn">) => useQuery<TData, TError>({ queryKey:
Common.UseTaskServiceGetTaskKeyFn({ dagId, taskId, versionNumber }, queryKey),
queryFn: () => TaskService.getTask({ dagId, taskId, versionNumber }) as TData,
...options });
/**
* Get Variable
* Get a variable entry.
diff --git a/airflow-core/src/airflow/ui/openapi-gen/queries/suspense.ts
b/airflow-core/src/airflow/ui/openapi-gen/queries/suspense.ts
index d51414ae7d7..6d3396e7bb2 100644
--- a/airflow-core/src/airflow/ui/openapi-gen/queries/suspense.ts
+++ b/airflow-core/src/airflow/ui/openapi-gen/queries/suspense.ts
@@ -1893,16 +1893,22 @@ export const useTaskServiceGetTasksSuspense = <TData =
Common.TaskServiceGetTask
/**
* Get Task
* Get simplified representation of a task.
+*
+* Pass ``version_number`` (e.g. a task instance's
``dag_version.version_number``) to get the task
+* as defined in that Dag version. Without it the latest version is used.
Returns 404 if that
+* version of the Dag does not exist.
* @param data The data for the request.
* @param data.dagId
* @param data.taskId
+* @param data.versionNumber
* @returns TaskResponse Successful Response
* @throws ApiError
*/
-export const useTaskServiceGetTaskSuspense = <TData =
Common.TaskServiceGetTaskDefaultResponse, TError = unknown, TQueryKey extends
Array<unknown> = unknown[]>({ dagId, taskId }: {
+export const useTaskServiceGetTaskSuspense = <TData =
Common.TaskServiceGetTaskDefaultResponse, TError = unknown, TQueryKey extends
Array<unknown> = unknown[]>({ dagId, taskId, versionNumber }: {
dagId: string;
taskId: unknown;
-}, queryKey?: TQueryKey, options?: Omit<UseQueryOptions<TData, TError>,
"queryKey" | "queryFn">) => useSuspenseQuery<TData, TError>({ queryKey:
Common.UseTaskServiceGetTaskKeyFn({ dagId, taskId }, queryKey), queryFn: () =>
TaskService.getTask({ dagId, taskId }) as TData, ...options });
+ versionNumber?: number;
+}, queryKey?: TQueryKey, options?: Omit<UseQueryOptions<TData, TError>,
"queryKey" | "queryFn">) => useSuspenseQuery<TData, TError>({ queryKey:
Common.UseTaskServiceGetTaskKeyFn({ dagId, taskId, versionNumber }, queryKey),
queryFn: () => TaskService.getTask({ dagId, taskId, versionNumber }) as TData,
...options });
/**
* Get Variable
* Get a variable entry.
diff --git a/airflow-core/src/airflow/ui/openapi-gen/requests/services.gen.ts
b/airflow-core/src/airflow/ui/openapi-gen/requests/services.gen.ts
index 1c89adc1f8d..8d0d3c6c05e 100644
--- a/airflow-core/src/airflow/ui/openapi-gen/requests/services.gen.ts
+++ b/airflow-core/src/airflow/ui/openapi-gen/requests/services.gen.ts
@@ -4635,9 +4635,14 @@ export class TaskService {
/**
* Get Task
* Get simplified representation of a task.
+ *
+ * Pass ``version_number`` (e.g. a task instance's
``dag_version.version_number``) to get the task
+ * as defined in that Dag version. Without it the latest version is used.
Returns 404 if that
+ * version of the Dag does not exist.
* @param data The data for the request.
* @param data.dagId
* @param data.taskId
+ * @param data.versionNumber
* @returns TaskResponse Successful Response
* @throws ApiError
*/
@@ -4649,6 +4654,9 @@ export class TaskService {
dag_id: data.dagId,
task_id: data.taskId
},
+ query: {
+ version_number: data.versionNumber
+ },
errors: {
400: 'Bad Request',
401: 'Unauthorized',
diff --git a/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts
b/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts
index 490a2c00d27..e27a29d5970 100644
--- a/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts
+++ b/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts
@@ -5033,6 +5033,7 @@ export type GetTasksResponse = TaskCollectionResponse;
export type GetTaskData = {
dagId: string;
taskId: unknown;
+ versionNumber?: number | null;
};
export type GetTaskResponse = TaskResponse;
diff --git
a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_tasks.py
b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_tasks.py
index 9fc48a9bbb8..3ba371ca079 100644
--- a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_tasks.py
+++ b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_tasks.py
@@ -21,6 +21,7 @@ from datetime import datetime
import pytest
from airflow.api_fastapi.common.dagbag import dag_bag_from_app
+from airflow.models.dag_version import DagVersion
from airflow.models.dagbag import DBDagBag
from airflow.providers.standard.operators.empty import EmptyOperator
from airflow.sdk import DAG
@@ -278,6 +279,40 @@ class TestGetTask(TestTaskEndpoint):
assert response.status_code == 200
assert response.json() == expected
+ def test_version_number_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"
+
+ with dag_maker(dag_id, session=session, bundle_version="commit-one"):
+ EmptyOperator(task_id=self.task_id, retries=4)
+ session.commit()
+ first_version_number = DagVersion.get_version(dag_id,
session=session).version_number
+
+ with dag_maker(dag_id, session=session, bundle_version="commit-two"):
+ EmptyOperator(task_id=self.task_id, retries=0)
+ session.commit()
+
+ assert DagVersion.get_version(dag_id, session=session).version_number
!= first_version_number
+
+ 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={"version_number":
first_version_number})
+ assert pinned.status_code == 200
+ assert pinned.json()["retries"] == 4
+
+ @pytest.mark.parametrize("version_number", [0, 99], ids=["zero",
"never-written"])
+ def test_unknown_version_number_is_not_found(self, test_client,
version_number):
+ """Without this, an unresolvable version would reach ``get_task`` and
raise a 500."""
+ response = test_client.get(
+ f"{self.api_prefix}/{self.dag_id}/tasks/{self.task_id}",
+ params={"version_number": version_number},
+ )
+ assert response.status_code == 404
+
def test_should_respond_404(self, test_client):
task_id = "xxxx_not_existing"
response = test_client.get(