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

1""" 

2Routes for interacting with flow objects. 

3""" 

4 

5from typing import List, Optional 

6from uuid import UUID 

7 

8from fastapi import Depends, HTTPException, Path, Response, status 

9from fastapi.param_functions import Body 

10 

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 

21 

22router: PrefectRouter = PrefectRouter(prefix="/flows", tags=["Flows"]) 

23 

24 

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. 

33 

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

38 

39 right_now = now("UTC") 

40 

41 async with db.session_context(begin_transaction=True) as session: 

42 model = await models.flows.create_flow(session=session, flow=flow) 

43 

44 if model.created >= right_now: 

45 response.status_code = status.HTTP_201_CREATED 

46 return model 

47 

48 

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 ) 

66 

67 

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 ) 

89 

90 

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 

106 

107 

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 

123 

124 

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 ) 

152 

153 

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 ) 

168 

169 

170BULK_OPERATION_LIMIT = 50 

171 

172 

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. 

188 

189 This also deletes all associated deployments. 

190 

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 ) 

200 

201 if not db_flows: 

202 return FlowBulkDeleteResponse(deleted=[]) 

203 

204 flow_ids = [f.id for f in db_flows] 

205 

206 # Delete flows (and their deployments) 

207 deleted_ids = await models.flows.delete_flows( 

208 session=session, 

209 flow_ids=flow_ids, 

210 ) 

211 

212 return FlowBulkDeleteResponse(deleted=deleted_ids) 

213 

214 

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 

231 

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 ) 

244 

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 ) 

253 

254 return FlowPaginationResponse( 

255 results=results, 

256 count=count, 

257 limit=limit, 

258 pages=(count + limit - 1) // limit, 

259 page=page, 

260 )