Coverage for /home/airflow/.local/lib/python3.12/site-packages/airflow/api_fastapi/core_api/routes/ui/gantt.py: 58%

24 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 fastapi import Depends, HTTPException, status 

21from sqlalchemy import or_, select, union_all 

22 

23from airflow.api_fastapi.auth.managers.models.resource_details import DagAccessEntity 

24from airflow.api_fastapi.common.db.common import SessionDep 

25from airflow.api_fastapi.common.router import AirflowRouter 

26from airflow.api_fastapi.core_api.datamodels.ui.gantt import GanttResponse, GanttTaskInstance 

27from airflow.api_fastapi.core_api.openapi.exceptions import create_openapi_http_exception_doc 

28from airflow.api_fastapi.core_api.security import requires_access_dag 

29from airflow.models.taskinstance import TaskInstance 

30from airflow.models.taskinstancehistory import TaskInstanceHistory 

31from airflow.utils.state import TaskInstanceState 

32 

33gantt_router = AirflowRouter(prefix="/gantt", tags=["Gantt"]) 

34 

35 

36@gantt_router.get( 

37 "/{dag_id}/{run_id}", 

38 responses=create_openapi_http_exception_doc( 

39 [ 

40 status.HTTP_404_NOT_FOUND, 

41 ] 

42 ), 

43 dependencies=[ 

44 Depends( 

45 requires_access_dag( 

46 method="GET", 

47 access_entity=DagAccessEntity.TASK_INSTANCE, 

48 ) 

49 ), 

50 Depends( 

51 requires_access_dag( 

52 method="GET", 

53 access_entity=DagAccessEntity.RUN, 

54 ) 

55 ), 

56 ], 

57) 

58def get_gantt_data( 

59 dag_id: str, 

60 run_id: str, 

61 session: SessionDep, 

62) -> GanttResponse: 

63 """Get all task instance tries for Gantt chart.""" 

64 # Exclude mapped tasks (use grid summaries) and UP_FOR_RETRY (already in history) 

65 current_tis = select( 

66 TaskInstance.task_id.label("task_id"), 

67 TaskInstance.task_display_name.label("task_display_name"), # type: ignore[attr-defined] 

68 TaskInstance.try_number.label("try_number"), 

69 TaskInstance.state.label("state"), 

70 TaskInstance.scheduled_dttm.label("scheduled_dttm"), 

71 TaskInstance.queued_dttm.label("queued_dttm"), 

72 TaskInstance.start_date.label("start_date"), 

73 TaskInstance.end_date.label("end_date"), 

74 ).where( 

75 TaskInstance.dag_id == dag_id, 

76 TaskInstance.run_id == run_id, 

77 TaskInstance.map_index == -1, 

78 or_(TaskInstance.state != TaskInstanceState.UP_FOR_RETRY, TaskInstance.state.is_(None)), 

79 ) 

80 

81 history_tis = select( 

82 TaskInstanceHistory.task_id.label("task_id"), 

83 TaskInstanceHistory.task_display_name.label("task_display_name"), 

84 TaskInstanceHistory.try_number.label("try_number"), 

85 TaskInstanceHistory.state.label("state"), 

86 TaskInstanceHistory.scheduled_dttm.label("scheduled_dttm"), 

87 TaskInstanceHistory.queued_dttm.label("queued_dttm"), 

88 TaskInstanceHistory.start_date.label("start_date"), 

89 TaskInstanceHistory.end_date.label("end_date"), 

90 ).where( 

91 TaskInstanceHistory.dag_id == dag_id, 

92 TaskInstanceHistory.run_id == run_id, 

93 TaskInstanceHistory.map_index == -1, 

94 ) 

95 

96 combined = union_all(current_tis, history_tis).subquery() 

97 query = select(combined).order_by(combined.c.task_id, combined.c.try_number) 

98 

99 results = session.execute(query).fetchall() 

100 

101 if not results: 

102 raise HTTPException( 

103 status.HTTP_404_NOT_FOUND, 

104 f"No task instances for dag_id={dag_id} run_id={run_id}", 

105 ) 

106 

107 task_instances = [ 

108 GanttTaskInstance( 

109 task_id=row.task_id, 

110 task_display_name=row.task_display_name, 

111 try_number=row.try_number, 

112 state=row.state, 

113 scheduled_dttm=row.scheduled_dttm, 

114 queued_dttm=row.queued_dttm, 

115 start_date=row.start_date, 

116 end_date=row.end_date, 

117 ) 

118 for row in results 

119 ] 

120 

121 return GanttResponse(dag_id=dag_id, run_id=run_id, task_instances=task_instances)