Coverage for /home/airflow/.local/lib/python3.12/site-packages/airflow/api_fastapi/core_api/routes/public/pools.py: 98%

51 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 Annotated, cast 

20 

21from fastapi import Depends, HTTPException, Query, status 

22from sqlalchemy import delete, select 

23from sqlalchemy.engine import CursorResult 

24 

25from airflow.api_fastapi.common.db.common import SessionDep, paginated_select 

26from airflow.api_fastapi.common.parameters import ( 

27 QueryLimit, 

28 QueryOffset, 

29 QueryPoolNamePatternSearch, 

30 QueryPoolNamePrefixPatternSearch, 

31 SortParam, 

32) 

33from airflow.api_fastapi.common.router import AirflowRouter 

34from airflow.api_fastapi.core_api.datamodels.common import BulkBody, BulkResponse 

35from airflow.api_fastapi.core_api.datamodels.pools import ( 

36 PoolBody, 

37 PoolCollectionResponse, 

38 PoolPatchBody, 

39 PoolResponse, 

40) 

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

42from airflow.api_fastapi.core_api.security import ( 

43 ReadablePoolsFilterDep, 

44 requires_access_pool, 

45 requires_access_pool_bulk, 

46) 

47from airflow.api_fastapi.core_api.services.public.pools import BulkPoolService, update_orm_from_pydantic 

48from airflow.api_fastapi.logging.decorators import action_logging 

49from airflow.models.pool import Pool 

50 

51pools_router = AirflowRouter(tags=["Pool"], prefix="/pools") 

52 

53 

54@pools_router.delete( 

55 "/{pool_name:path}", 

56 status_code=status.HTTP_204_NO_CONTENT, 

57 responses=create_openapi_http_exception_doc( 

58 [ 

59 status.HTTP_400_BAD_REQUEST, 

60 status.HTTP_404_NOT_FOUND, 

61 ] 

62 ), 

63 dependencies=[Depends(requires_access_pool(method="DELETE")), Depends(action_logging())], 

64) 

65def delete_pool( 

66 pool_name: str, 

67 session: SessionDep, 

68): 

69 """Delete a pool entry.""" 

70 if pool_name == "default_pool": 

71 raise HTTPException(status.HTTP_400_BAD_REQUEST, "Default Pool can't be deleted") 

72 

73 affected_count = cast("CursorResult", session.execute(delete(Pool).where(Pool.pool == pool_name))) 

74 

75 if affected_count.rowcount == 0: 

76 raise HTTPException(status.HTTP_404_NOT_FOUND, f"The Pool with name: `{pool_name}` was not found") 

77 

78 

79@pools_router.get( 

80 "/{pool_name:path}", 

81 responses=create_openapi_http_exception_doc([status.HTTP_404_NOT_FOUND]), 

82 dependencies=[Depends(requires_access_pool(method="GET"))], 

83) 

84def get_pool( 

85 pool_name: str, 

86 session: SessionDep, 

87) -> PoolResponse: 

88 """Get a pool.""" 

89 pool = session.scalar(select(Pool).where(Pool.pool == pool_name)) 

90 if pool is None: 

91 raise HTTPException(status.HTTP_404_NOT_FOUND, f"The Pool with name: `{pool_name}` was not found") 

92 

93 return pool 

94 

95 

96@pools_router.get( 

97 "", 

98 dependencies=[Depends(requires_access_pool(method="GET"))], 

99) 

100def get_pools( 

101 limit: QueryLimit, 

102 offset: QueryOffset, 

103 order_by: Annotated[ 

104 SortParam, 

105 Depends(SortParam(["id", "pool"], Pool, to_replace={"name": "pool"}).dynamic_depends()), 

106 ], 

107 pool_name_pattern: QueryPoolNamePatternSearch, 

108 pool_name_prefix_pattern: QueryPoolNamePrefixPatternSearch, 

109 readable_pools_filter: ReadablePoolsFilterDep, 

110 session: SessionDep, 

111) -> PoolCollectionResponse: 

112 """Get all pools entries.""" 

113 pools_select, total_entries = paginated_select( 

114 statement=select(Pool), 

115 filters=[pool_name_pattern, pool_name_prefix_pattern, readable_pools_filter], 

116 order_by=order_by, 

117 offset=offset, 

118 limit=limit, 

119 session=session, 

120 ) 

121 

122 pools = session.scalars(pools_select) 

123 

124 return PoolCollectionResponse( 

125 pools=pools, 

126 total_entries=total_entries, 

127 ) 

128 

129 

130@pools_router.patch( 

131 "/{pool_name:path}", 

132 responses=create_openapi_http_exception_doc( 

133 [ 

134 status.HTTP_400_BAD_REQUEST, 

135 status.HTTP_404_NOT_FOUND, 

136 ] 

137 ), 

138 dependencies=[Depends(requires_access_pool(method="PUT")), Depends(action_logging())], 

139) 

140def patch_pool( 

141 pool_name: str, 

142 patch_body: PoolPatchBody, 

143 session: SessionDep, 

144 update_mask: list[str] | None = Query(None), 

145) -> PoolResponse: 

146 """Update a Pool.""" 

147 if patch_body.name and patch_body.name != pool_name: 

148 raise HTTPException( 

149 status.HTTP_400_BAD_REQUEST, 

150 "Invalid body, pool name from request body doesn't match uri parameter", 

151 ) 

152 

153 pool = session.scalar(select(Pool).where(Pool.pool == pool_name).limit(1)) 

154 if not pool: 

155 raise HTTPException( 

156 status.HTTP_404_NOT_FOUND, detail=f"The Pool with name: `{pool_name}` was not found" 

157 ) 

158 

159 updated_pool = update_orm_from_pydantic(pool, patch_body, update_mask) 

160 return updated_pool 

161 

162 

163@pools_router.post( 

164 "", 

165 status_code=status.HTTP_201_CREATED, 

166 responses=create_openapi_http_exception_doc( 

167 [status.HTTP_409_CONFLICT] 

168 ), # handled by global exception handler 

169 dependencies=[Depends(requires_access_pool(method="POST")), Depends(action_logging())], 

170) 

171def post_pool( 

172 body: PoolBody, 

173 session: SessionDep, 

174) -> PoolResponse: 

175 """Create a Pool.""" 

176 pool = Pool(**body.model_dump()) 

177 session.add(pool) 

178 return pool 

179 

180 

181@pools_router.patch( 

182 "", 

183 dependencies=[Depends(requires_access_pool_bulk()), Depends(action_logging())], 

184) 

185def bulk_pools( 

186 request: BulkBody[PoolBody], 

187 session: SessionDep, 

188) -> BulkResponse: 

189 """Bulk create, update, and delete pools.""" 

190 return BulkPoolService(session=session, request=request).handle_request()