Coverage for /home/airflow/.local/lib/python3.12/site-packages/airflow/api_fastapi/core_api/routes/ui/dag_runs.py: 71%
19 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 select
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.dag_runs import DagRunStatsResponse
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.api_fastapi.core_api.services.ui.dag_run import compute_duration_stats
30from airflow.models.dagrun import DagRun
31from airflow.utils.state import DagRunState
33dag_runs_router = AirflowRouter(prefix="/dags/{dag_id}/dagRuns", tags=["DagRun"])
36@dag_runs_router.get(
37 "/{dag_run_id}/stats",
38 responses=create_openapi_http_exception_doc([status.HTTP_404_NOT_FOUND]),
39 dependencies=[Depends(requires_access_dag(method="GET", access_entity=DagAccessEntity.RUN))],
40)
41def get_dag_run_stats(dag_id: str, dag_run_id: str, session: SessionDep) -> DagRunStatsResponse:
42 """Get duration statistics for a DAG based on its historical completed runs."""
43 if not session.scalar(select(DagRun.id).filter_by(dag_id=dag_id, run_id=dag_run_id)):
44 raise HTTPException(
45 status.HTTP_404_NOT_FOUND,
46 f"The DagRun with dag_id: `{dag_id}` and run_id: `{dag_run_id}` was not found",
47 )
49 durations = [
50 d
51 for d in session.scalars(
52 select(DagRun.duration.expression) # type: ignore[attr-defined]
53 .where(
54 DagRun.dag_id == dag_id,
55 DagRun.state.in_([DagRunState.SUCCESS, DagRunState.FAILED]),
56 )
57 .order_by(DagRun.run_after.desc())
58 .limit(100)
59 )
60 if d is not None
61 ]
63 return DagRunStatsResponse(duration=compute_duration_stats(durations))