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
« 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 Annotated, cast
21from fastapi import Depends, HTTPException, Query, status
22from sqlalchemy import delete, select
23from sqlalchemy.engine import CursorResult
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
51pools_router = AirflowRouter(tags=["Pool"], prefix="/pools")
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")
73 affected_count = cast("CursorResult", session.execute(delete(Pool).where(Pool.pool == pool_name)))
75 if affected_count.rowcount == 0:
76 raise HTTPException(status.HTTP_404_NOT_FOUND, f"The Pool with name: `{pool_name}` was not found")
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")
93 return pool
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 )
122 pools = session.scalars(pools_select)
124 return PoolCollectionResponse(
125 pools=pools,
126 total_entries=total_entries,
127 )
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 )
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 )
159 updated_pool = update_orm_from_pydantic(pool, patch_body, update_mask)
160 return updated_pool
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
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()