Coverage for /usr/local/lib/python3.12/site-packages/prefect/server/api/flows.py: 92%
72 statements
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-07 02:04 +0000
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-07 02:04 +0000
1"""
2Routes for interacting with flow objects.
3"""
5from typing import List, Optional
6from uuid import UUID
8from fastapi import Depends, HTTPException, Path, Response, status
9from fastapi.param_functions import Body
11import prefect.server.api.dependencies as dependencies
12import prefect.server.models as models
13import prefect.server.schemas as schemas
14from prefect.server.database import PrefectDBInterface, provide_database_interface
15from prefect.server.schemas.responses import (
16 FlowBulkDeleteResponse,
17 FlowPaginationResponse,
18)
19from prefect.server.utilities.server import PrefectRouter
20from prefect.types._datetime import now
22router: PrefectRouter = PrefectRouter(prefix="/flows", tags=["Flows"])
25@router.post("/")
26async def create_flow(
27 flow: schemas.actions.FlowCreate,
28 response: Response,
29 db: PrefectDBInterface = Depends(provide_database_interface),
30) -> schemas.core.Flow:
31 """Creates a new flow from the provided schema. If a flow with the
32 same name already exists, the existing flow is returned.
34 For more information, see https://docs.prefect.io/v3/concepts/flows.
35 """
36 # hydrate the input model into a full flow model
37 flow = schemas.core.Flow(**flow.model_dump())
39 right_now = now("UTC")
41 async with db.session_context(begin_transaction=True) as session:
42 model = await models.flows.create_flow(session=session, flow=flow)
44 if model.created >= right_now:
45 response.status_code = status.HTTP_201_CREATED
46 return model
49@router.patch("/{id:uuid}", status_code=status.HTTP_204_NO_CONTENT)
50async def update_flow(
51 flow: schemas.actions.FlowUpdate,
52 flow_id: UUID = Path(..., description="The flow id", alias="id"),
53 db: PrefectDBInterface = Depends(provide_database_interface),
54) -> None:
55 """
56 Updates a flow.
57 """
58 async with db.session_context(begin_transaction=True) as session:
59 result = await models.flows.update_flow(
60 session=session, flow=flow, flow_id=flow_id
61 )
62 if not result:
63 raise HTTPException(
64 status_code=status.HTTP_404_NOT_FOUND, detail="Flow not found"
65 )
68@router.post("/count")
69async def count_flows(
70 flows: schemas.filters.FlowFilter = None,
71 flow_runs: schemas.filters.FlowRunFilter = None,
72 task_runs: schemas.filters.TaskRunFilter = None,
73 deployments: schemas.filters.DeploymentFilter = None,
74 work_pools: schemas.filters.WorkPoolFilter = None,
75 db: PrefectDBInterface = Depends(provide_database_interface),
76) -> int:
77 """
78 Count flows.
79 """
80 async with db.session_context() as session:
81 return await models.flows.count_flows(
82 session=session,
83 flow_filter=flows,
84 flow_run_filter=flow_runs,
85 task_run_filter=task_runs,
86 deployment_filter=deployments,
87 work_pool_filter=work_pools,
88 )
91@router.get("/name/{name}")
92async def read_flow_by_name(
93 name: str = Path(..., description="The name of the flow"),
94 db: PrefectDBInterface = Depends(provide_database_interface),
95) -> schemas.core.Flow:
96 """
97 Get a flow by name.
98 """
99 async with db.session_context() as session:
100 flow = await models.flows.read_flow_by_name(session=session, name=name)
101 if not flow: 101 ↛ 105line 101 didn't jump to line 105 because the condition on line 101 was always true
102 raise HTTPException(
103 status_code=status.HTTP_404_NOT_FOUND, detail="Flow not found"
104 )
105 return flow
108@router.get("/{id:uuid}")
109async def read_flow(
110 flow_id: UUID = Path(..., description="The flow id", alias="id"),
111 db: PrefectDBInterface = Depends(provide_database_interface),
112) -> schemas.core.Flow:
113 """
114 Get a flow by id.
115 """
116 async with db.session_context() as session:
117 flow = await models.flows.read_flow(session=session, flow_id=flow_id)
118 if not flow:
119 raise HTTPException(
120 status_code=status.HTTP_404_NOT_FOUND, detail="Flow not found"
121 )
122 return flow
125@router.post("/filter")
126async def read_flows(
127 limit: int = dependencies.LimitBody(),
128 offset: int = Body(0, ge=0),
129 flows: schemas.filters.FlowFilter = None,
130 flow_runs: schemas.filters.FlowRunFilter = None,
131 task_runs: schemas.filters.TaskRunFilter = None,
132 deployments: schemas.filters.DeploymentFilter = None,
133 work_pools: schemas.filters.WorkPoolFilter = None,
134 sort: schemas.sorting.FlowSort = Body(schemas.sorting.FlowSort.NAME_ASC),
135 db: PrefectDBInterface = Depends(provide_database_interface),
136) -> List[schemas.core.Flow]:
137 """
138 Query for flows.
139 """
140 async with db.session_context() as session:
141 return await models.flows.read_flows(
142 session=session,
143 flow_filter=flows,
144 flow_run_filter=flow_runs,
145 task_run_filter=task_runs,
146 deployment_filter=deployments,
147 work_pool_filter=work_pools,
148 sort=sort,
149 offset=offset,
150 limit=limit,
151 )
154@router.delete("/{id:uuid}", status_code=status.HTTP_204_NO_CONTENT)
155async def delete_flow(
156 flow_id: UUID = Path(..., description="The flow id", alias="id"),
157 db: PrefectDBInterface = Depends(provide_database_interface),
158) -> None:
159 """
160 Delete a flow by id.
161 """
162 async with db.session_context(begin_transaction=True) as session:
163 result = await models.flows.delete_flow(session=session, flow_id=flow_id)
164 if not result:
165 raise HTTPException(
166 status_code=status.HTTP_404_NOT_FOUND, detail="Flow not found"
167 )
170BULK_OPERATION_LIMIT = 50
173@router.post("/bulk_delete")
174async def bulk_delete_flows(
175 flows: Optional[schemas.filters.FlowFilter] = Body(
176 None, description="Filter criteria for flows to delete"
177 ),
178 limit: int = Body(
179 BULK_OPERATION_LIMIT,
180 ge=1,
181 le=BULK_OPERATION_LIMIT,
182 description=f"Maximum number of flows to delete. Defaults to {BULK_OPERATION_LIMIT}.",
183 ),
184 db: PrefectDBInterface = Depends(provide_database_interface),
185) -> FlowBulkDeleteResponse:
186 """
187 Bulk delete flows matching the specified filter criteria.
189 This also deletes all associated deployments.
191 Returns the IDs of flows that were deleted.
192 """
193 async with db.session_context(begin_transaction=True) as session:
194 # Query matching flows
195 db_flows = await models.flows.read_flows(
196 session=session,
197 flow_filter=flows,
198 limit=limit,
199 )
201 if not db_flows:
202 return FlowBulkDeleteResponse(deleted=[])
204 flow_ids = [f.id for f in db_flows]
206 # Delete flows (and their deployments)
207 deleted_ids = await models.flows.delete_flows(
208 session=session,
209 flow_ids=flow_ids,
210 )
212 return FlowBulkDeleteResponse(deleted=deleted_ids)
215@router.post("/paginate")
216async def paginate_flows(
217 limit: int = dependencies.LimitBody(),
218 page: int = Body(1, ge=1),
219 flows: Optional[schemas.filters.FlowFilter] = None,
220 flow_runs: Optional[schemas.filters.FlowRunFilter] = None,
221 task_runs: Optional[schemas.filters.TaskRunFilter] = None,
222 deployments: Optional[schemas.filters.DeploymentFilter] = None,
223 work_pools: Optional[schemas.filters.WorkPoolFilter] = None,
224 sort: schemas.sorting.FlowSort = Body(schemas.sorting.FlowSort.NAME_ASC),
225 db: PrefectDBInterface = Depends(provide_database_interface),
226) -> FlowPaginationResponse:
227 """
228 Pagination query for flows.
229 """
230 offset = (page - 1) * limit
232 async with db.session_context() as session:
233 results = await models.flows.read_flows(
234 session=session,
235 flow_filter=flows,
236 flow_run_filter=flow_runs,
237 task_run_filter=task_runs,
238 deployment_filter=deployments,
239 work_pool_filter=work_pools,
240 sort=sort,
241 offset=offset,
242 limit=limit,
243 )
245 count = await models.flows.count_flows(
246 session=session,
247 flow_filter=flows,
248 flow_run_filter=flow_runs,
249 task_run_filter=task_runs,
250 deployment_filter=deployments,
251 work_pool_filter=work_pools,
252 )
254 return FlowPaginationResponse(
255 results=results,
256 count=count,
257 limit=limit,
258 pages=(count + limit - 1) // limit,
259 page=page,
260 )