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
« 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 fastapi import Depends, HTTPException, status
21from sqlalchemy import or_, select, union_all
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
33gantt_router = AirflowRouter(prefix="/gantt", tags=["Gantt"])
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 )
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 )
96 combined = union_all(current_tis, history_tis).subquery()
97 query = select(combined).order_by(combined.c.task_id, combined.c.try_number)
99 results = session.execute(query).fetchall()
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 )
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 ]
121 return GanttResponse(dag_id=dag_id, run_id=run_id, task_instances=task_instances)