Coverage for /home/airflow/.local/lib/python3.12/site-packages/airflow/api_fastapi/core_api/routes/public/tasks.py: 100%
29 statements
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-07 14:22 +0000
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-07 14:22 +0000
1# Licensed to the Apache Software Foundation (ASF) under one
2# or more contributor license agreements. See the NOTICE file
3# distributed with this work for additional information
4# regarding copyright ownership. The ASF licenses this file
5# to you under the Apache License, Version 2.0 (the
6# "License"); you may not use this file except in compliance
7# with the License. You may obtain a copy of the License at
8#
9# http://www.apache.org/licenses/LICENSE-2.0
10#
11# Unless required by applicable law or agreed to in writing,
12# software distributed under the License is distributed on an
13# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
14# KIND, either express or implied. See the License for the
15# specific language governing permissions and limitations
16# under the License.
18from __future__ import annotations
20from typing import cast
22from fastapi import Depends, HTTPException, status
24from airflow.api_fastapi.auth.managers.models.resource_details import DagAccessEntity
25from airflow.api_fastapi.common.dagbag import DagBagDep, get_latest_version_of_dag
26from airflow.api_fastapi.common.db.common import SessionDep
27from airflow.api_fastapi.common.router import AirflowRouter
28from airflow.api_fastapi.core_api.datamodels.tasks import TaskCollectionResponse, TaskResponse
29from airflow.api_fastapi.core_api.openapi.exceptions import create_openapi_http_exception_doc
30from airflow.api_fastapi.core_api.security import requires_access_dag
31from airflow.exceptions import TaskNotFound
33tasks_router = AirflowRouter(tags=["Task"], prefix="/dags/{dag_id}/tasks")
35_SORTABLE_TASK_FIELDS = {
36 "task_id",
37 "task_display_name",
38 "owner",
39 "start_date",
40 "end_date",
41 "trigger_rule",
42 "depends_on_past",
43 "wait_for_downstream",
44 "retries",
45 "queue",
46 "pool",
47 "pool_slots",
48 "execution_timeout",
49 "retry_delay",
50 "retry_exponential_backoff",
51 "priority_weight",
52 "weight_rule",
53 "ui_color",
54 "ui_fgcolor",
55 "operator_name",
56}
59@tasks_router.get(
60 "",
61 responses=create_openapi_http_exception_doc(
62 [
63 status.HTTP_400_BAD_REQUEST,
64 status.HTTP_404_NOT_FOUND,
65 ]
66 ),
67 dependencies=[Depends(requires_access_dag(method="GET", access_entity=DagAccessEntity.TASK))],
68)
69def get_tasks(
70 dag_id: str,
71 dag_bag: DagBagDep,
72 session: SessionDep,
73 order_by: str = "task_id",
74) -> TaskCollectionResponse:
75 """Get tasks for Dag."""
76 dag = get_latest_version_of_dag(dag_bag, dag_id, session)
77 lstripped_order_by = order_by.lstrip("-")
78 if lstripped_order_by not in _SORTABLE_TASK_FIELDS:
79 raise HTTPException(
80 status.HTTP_400_BAD_REQUEST,
81 f"Ordering with '{lstripped_order_by}' is disallowed or "
82 f"the attribute does not exist on the model",
83 )
84 tasks = sorted(
85 dag.tasks,
86 key=lambda task: (getattr(task, lstripped_order_by) is None, getattr(task, lstripped_order_by)),
87 reverse=(order_by[0:1] == "-"),
88 )
89 return TaskCollectionResponse(
90 tasks=cast("list[TaskResponse]", tasks),
91 total_entries=len(tasks),
92 )
95@tasks_router.get(
96 "/{task_id}",
97 responses=create_openapi_http_exception_doc(
98 [
99 status.HTTP_400_BAD_REQUEST,
100 status.HTTP_404_NOT_FOUND,
101 ]
102 ),
103 dependencies=[Depends(requires_access_dag(method="GET", access_entity=DagAccessEntity.TASK))],
104)
105def get_task(dag_id: str, task_id, session: SessionDep, dag_bag: DagBagDep) -> TaskResponse:
106 """Get simplified representation of a task."""
107 dag = get_latest_version_of_dag(dag_bag, dag_id, session)
108 try:
109 task = dag.get_task(task_id=task_id)
110 except TaskNotFound:
111 raise HTTPException(status.HTTP_404_NOT_FOUND, f"Task with id {task_id} was not found")
112 return cast("TaskResponse", task)