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

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. 

17 

18from __future__ import annotations 

19 

20from typing import cast 

21 

22from fastapi import Depends, HTTPException, status 

23 

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 

32 

33tasks_router = AirflowRouter(tags=["Task"], prefix="/dags/{dag_id}/tasks") 

34 

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} 

57 

58 

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 ) 

93 

94 

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)