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

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 

18 

19from typing import TYPE_CHECKING, Annotated 

20 

21from fastapi import Depends, HTTPException, Request, status 

22from sqlalchemy.orm import Session 

23 

24from airflow.configuration import conf 

25from airflow.models.dagbag import CachedDBDagBag, DBDagBag 

26from airflow.models.serialized_dag import SerializedDagModel 

27 

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 

31 

32 

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) 

37 

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") 

42 

43 return CachedDBDagBag( 

44 cache_size=cache_size, 

45 cache_ttl=cache_ttl, 

46 stats_prefix="api_server.dag_bag", 

47 ) 

48 

49 

50def dag_bag_from_app(request: Request) -> DBDagBag: 

51 """ 

52 FastAPI dependency resolver that returns the shared DagBag instance from app.state. 

53 

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 

58 

59 

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 

75 

76 

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 

82 

83 

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. 

89 

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. 

93 

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 

109 

110 

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. 

119 

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) 

133 

134 

135DagBagDep = Annotated[DBDagBag, Depends(dag_bag_from_app)]