Coverage for /home/airflow/.local/lib/python3.12/site-packages/airflow/api_fastapi/common/dagbag.py: 78%
52 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.
17from __future__ import annotations
19from typing import TYPE_CHECKING, Annotated
21from fastapi import Depends, HTTPException, Request, status
22from sqlalchemy.orm import Session
24from airflow.configuration import conf
25from airflow.models.dagbag import CachedDBDagBag, DBDagBag
26from airflow.models.serialized_dag import SerializedDagModel
28if TYPE_CHECKING: 28 ↛ 29line 28 didn't jump to line 29 because the condition on line 28 was never true
29 from airflow.models.dagrun import DagRun
30 from airflow.serialization.definitions.dag import SerializedDAG
33def create_dag_bag() -> CachedDBDagBag:
34 """Create DagBag with configurable LRU+TTL caching for API server usage."""
35 cache_size = conf.getint("api", "dag_cache_size", fallback=64)
36 cache_ttl = conf.getint("api", "dag_cache_ttl", fallback=3600)
38 if cache_size < 0: 38 ↛ 39line 38 didn't jump to line 39 because the condition on line 38 was never true
39 raise ValueError("[api] dag_cache_size must be greater than or equal to 0")
40 if cache_ttl < 0: 40 ↛ 41line 40 didn't jump to line 41 because the condition on line 40 was never true
41 raise ValueError("[api] dag_cache_ttl must be greater than or equal to 0")
43 return CachedDBDagBag(
44 cache_size=cache_size,
45 cache_ttl=cache_ttl,
46 stats_prefix="api_server.dag_bag",
47 )
50def dag_bag_from_app(request: Request) -> DBDagBag:
51 """
52 FastAPI dependency resolver that returns the shared DagBag instance from app.state.
54 This ensures that all API routes using DagBag via dependency injection receive the same
55 singleton instance that was initialized at app startup.
56 """
57 return request.app.state.dag_bag
60def get_latest_version_of_dag(
61 dag_bag: DBDagBag, dag_id: str, session: Session, include_reason: bool = False
62) -> SerializedDAG:
63 dag = dag_bag.get_latest_version_of_dag(dag_id, session=session)
64 if not dag:
65 if include_reason: 65 ↛ 66line 65 didn't jump to line 66 because the condition on line 65 was never true
66 raise HTTPException(
67 status.HTTP_404_NOT_FOUND,
68 detail={
69 "reason": "not_found",
70 "message": f"The Dag with ID: `{dag_id}` was not found",
71 },
72 )
73 raise HTTPException(status.HTTP_404_NOT_FOUND, f"The Dag with ID: `{dag_id}` was not found")
74 return dag
77def get_dag_for_run(dag_bag: DBDagBag, dag_run: DagRun, session: Session) -> SerializedDAG:
78 dag = dag_bag.get_dag_for_run(dag_run, session=session)
79 if not dag: 79 ↛ 80line 79 didn't jump to line 80 because the condition on line 79 was never true
80 raise HTTPException(status.HTTP_404_NOT_FOUND, f"The Dag with ID: `{dag_run.dag_id}` was not found")
81 return dag
84def get_dag_for_run_or_latest_version(
85 dag_bag: DBDagBag, dag_run: DagRun | None, dag_id: str | None, session: Session
86) -> SerializedDAG:
87 """
88 Retrieve the serialized Dag for a specific run, or the latest version if no run is given.
90 When a dag_run is provided, we prefer the exact Dag version the run was created with
91 (``created_dag_version_id``) so that task group lookups, operator metadata, etc. match
92 the Dag structure at the time of the run.
94 This is necessary because ``get_dag_for_run`` delegates to ``_version_from_dag_run``
95 which, for unversioned bundles (e.g. ``LocalDagBundle``), falls back to the *latest*
96 ``DagVersion``.
97 """
98 dag: SerializedDAG | None = None
99 if dag_run:
100 if dag_run.created_dag_version_id: 100 ↛ 102line 100 didn't jump to line 102 because the condition on line 100 was always true
101 dag = dag_bag.get_dag(dag_run.created_dag_version_id, session=session)
102 if not dag: 102 ↛ 103line 102 didn't jump to line 103 because the condition on line 102 was never true
103 dag = dag_bag.get_dag_for_run(dag_run, session=session)
104 elif dag_id: 104 ↛ 106line 104 didn't jump to line 106 because the condition on line 104 was always true
105 dag = dag_bag.get_latest_version_of_dag(dag_id, session=session)
106 if not dag:
107 raise HTTPException(status.HTTP_404_NOT_FOUND, f"The Dag with ID: `{dag_id}` was not found")
108 return dag
111def resolve_run_on_latest_version(
112 explicit_value: bool | None,
113 dag_id: str,
114 session: Session,
115 fallback: bool = False,
116) -> bool:
117 """
118 Resolve run_on_latest_version using precedence: explicit > DAG-level > global config > fallback.
120 :param explicit_value: Value from the API request body (or None if not specified).
121 :param dag_id: The DAG ID to look up.
122 :param session: Database session.
123 :param fallback: Default to use when neither DAG-level nor global config is set.
124 Clear/rerun endpoints use False (the historical default).
125 Backfill endpoint uses True (the historical default for backfills).
126 """
127 if explicit_value is not None:
128 return explicit_value
129 serialized = SerializedDagModel.get_dag(dag_id, session=session)
130 if serialized and serialized.rerun_with_latest_version is not None: 130 ↛ 131line 130 didn't jump to line 131 because the condition on line 130 was never true
131 return serialized.rerun_with_latest_version
132 return conf.getboolean("core", "rerun_with_latest_version", fallback=fallback)
135DagBagDep = Annotated[DBDagBag, Depends(dag_bag_from_app)]