Coverage for /usr/local/lib/python3.12/site-packages/prefect/server/api/flow_runs.py: 58%

313 statements  

« prev     ^ index     » next       coverage.py v7.15.2, created at 2026-10-07 02:04 +0000

1""" 

2Routes for interacting with flow run objects. 

3""" 

4 

5import asyncio 

6import csv 

7import datetime 

8import io 

9from typing import TYPE_CHECKING, Any, Dict, List, Optional 

10from urllib.parse import quote 

11from uuid import UUID 

12 

13import orjson 

14import sqlalchemy as sa 

15from docket import Depends as DocketDepends 

16from docket import Docket, Retry 

17from fastapi import ( 

18 Body, 

19 Depends, 

20 HTTPException, 

21 Path, 

22 Query, 

23 Request, 

24 Response, 

25) 

26from fastapi.encoders import jsonable_encoder 

27from fastapi.responses import PlainTextResponse, StreamingResponse 

28from sqlalchemy.exc import IntegrityError 

29 

30import prefect.server.api.dependencies as dependencies 

31import prefect.server.models as models 

32import prefect.server.schemas as schemas 

33from prefect._internal.compatibility.starlette import status 

34from prefect.logging import get_logger 

35from prefect.server.api.run_history import run_history 

36from prefect.server.api.validation import validate_job_variables_for_deployment_flow_run 

37from prefect.server.api.workers import WorkerLookups 

38from prefect.server.database import PrefectDBInterface, provide_database_interface 

39from prefect.server.exceptions import FlowRunGraphTooLarge 

40from prefect.server.models.flow_runs import ( 

41 DependencyResult, 

42 read_flow_run_graph, 

43) 

44from prefect.server.orchestration import dependencies as orchestration_dependencies 

45from prefect.server.orchestration.policies import ( 

46 FlowRunOrchestrationPolicy, 

47 TaskRunOrchestrationPolicy, 

48) 

49from prefect.server.schemas.graph import Graph 

50from prefect.server.schemas.responses import ( 

51 FlowRunBulkDeleteResponse, 

52 FlowRunBulkSetStateResponse, 

53 FlowRunOrchestrationResult, 

54 FlowRunPaginationResponse, 

55 OrchestrationResult, 

56) 

57from prefect.server.services.cancellation_cleanup import ( 

58 maybe_schedule_cancelling_timeout_check_for_state, 

59) 

60from prefect.server.utilities.server import PrefectRouter 

61from prefect.types import DateTime 

62from prefect.types._datetime import earliest_possible_datetime, now 

63from prefect.utilities import schema_tools 

64 

65if TYPE_CHECKING: 65 ↛ 66line 65 didn't jump to line 66 because the condition on line 65 was never true

66 import logging 

67 

68logger: "logging.Logger" = get_logger("server.api") 

69 

70router: PrefectRouter = PrefectRouter(prefix="/flow_runs", tags=["Flow Runs"]) 

71 

72 

73def _get_request_docket(request: Request) -> Docket | None: 

74 return getattr(request.app.state, "docket", None) 

75 

76 

77async def _maybe_schedule_cancelling_timeout_check_for_state( 

78 *, 

79 request: Request, 

80 flow_run_id: UUID, 

81 state: schemas.states.State | None, 

82) -> None: 

83 docket = _get_request_docket(request) 

84 if docket is None: 84 ↛ 85line 84 didn't jump to line 85 because the condition on line 84 was never true

85 return 

86 

87 try: 

88 await maybe_schedule_cancelling_timeout_check_for_state( 

89 docket=docket, 

90 flow_run_id=flow_run_id, 

91 state=state, 

92 ) 

93 except Exception: 

94 logger.exception( 

95 "Failed to schedule CANCELLING timeout check; allowing accepted " 

96 "state transition to proceed", 

97 extra={ 

98 "flow_run_id": str(flow_run_id), 

99 "flow_run_state_id": str(state.id) if state and state.id else None, 

100 }, 

101 ) 

102 

103 

104@router.post("/") 

105async def create_flow_run( 

106 flow_run: schemas.actions.FlowRunCreate, 

107 request: Request, 

108 db: PrefectDBInterface = Depends(provide_database_interface), 

109 response: Response = None, # type: ignore 

110 created_by: Optional[schemas.core.CreatedBy] = Depends(dependencies.get_created_by), 

111 orchestration_parameters: Dict[str, Any] = Depends( 

112 orchestration_dependencies.provide_flow_orchestration_parameters 

113 ), 

114 api_version: str = Depends(dependencies.provide_request_api_version), 

115 worker_lookups: WorkerLookups = Depends(WorkerLookups), 

116) -> schemas.responses.FlowRunResponse: 

117 """ 

118 Create a flow run. If a flow run with the same flow_id and 

119 idempotency key already exists, the existing flow run will be returned. 

120 

121 If no state is provided, the flow run will be created in a PENDING state. 

122 

123 For more information, see https://docs.prefect.io/v3/concepts/flows. 

124 """ 

125 # hydrate the input model into a full flow run / state model 

126 flow_run_object = schemas.core.FlowRun( 

127 **flow_run.model_dump(), created_by=created_by 

128 ) 

129 

130 # pass the request version to the orchestration engine to support compatibility code 

131 orchestration_parameters.update({"api-version": api_version}) 

132 

133 if not flow_run_object.state: 

134 flow_run_object.state = schemas.states.Pending() 

135 

136 right_now = now("UTC") 

137 

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

139 if flow_run.work_pool_name: 

140 if flow_run.work_queue_name: 140 ↛ 141line 140 didn't jump to line 141 because the condition on line 140 was never true

141 work_queue_id = await worker_lookups._get_work_queue_id_from_name( 

142 session=session, 

143 work_pool_name=flow_run.work_pool_name, 

144 work_queue_name=flow_run.work_queue_name, 

145 ) 

146 else: 

147 work_queue_id = ( 

148 await worker_lookups._get_default_work_queue_id_from_work_pool_name( 

149 session=session, 

150 work_pool_name=flow_run.work_pool_name, 

151 ) 

152 ) 

153 else: 

154 work_queue_id = None 

155 

156 flow_run_object.work_queue_id = work_queue_id 

157 

158 model = await models.flow_runs.create_flow_run( 

159 session=session, 

160 flow_run=flow_run_object, 

161 orchestration_parameters=orchestration_parameters, 

162 ) 

163 created = model.created >= right_now 

164 timeout_check_state = ( 

165 schemas.states.State.from_orm_without_result(model.state) 

166 if created and model.state 

167 else None 

168 ) 

169 flow_run_id = model.id 

170 flow_run_response = schemas.responses.FlowRunResponse.model_validate( 

171 model, from_attributes=True 

172 ) 

173 

174 await _maybe_schedule_cancelling_timeout_check_for_state( 

175 request=request, 

176 flow_run_id=flow_run_id, 

177 state=timeout_check_state, 

178 ) 

179 if created: 179 ↛ 182line 179 didn't jump to line 182 because the condition on line 179 was always true

180 response.status_code = status.HTTP_201_CREATED 

181 

182 return flow_run_response 

183 

184 

185@router.patch("/{id:uuid}", status_code=status.HTTP_204_NO_CONTENT) 

186async def update_flow_run( 

187 flow_run: schemas.actions.FlowRunUpdate, 

188 flow_run_id: UUID = Path(..., description="The flow run id", alias="id"), 

189 db: PrefectDBInterface = Depends(provide_database_interface), 

190) -> None: 

191 """ 

192 Updates a flow run. 

193 """ 

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

195 if flow_run.job_variables is not None: 

196 this_run = await models.flow_runs.read_flow_run( 

197 session, flow_run_id=flow_run_id 

198 ) 

199 if this_run is None: 

200 raise HTTPException( 

201 status.HTTP_404_NOT_FOUND, detail="Flow run not found" 

202 ) 

203 if not this_run.state: 203 ↛ 204line 203 didn't jump to line 204 because the condition on line 203 was never true

204 raise HTTPException( 

205 status.HTTP_400_BAD_REQUEST, 

206 detail="Flow run state is required to update job variables but none exists", 

207 ) 

208 if this_run.state.type != schemas.states.StateType.SCHEDULED: 208 ↛ anywhereline 208 didn't jump anywhere: it always raised an exception.

209 raise HTTPException( 

210 status_code=status.HTTP_400_BAD_REQUEST, 

211 detail=f"Job variables for a flow run in state {this_run.state.type.name} cannot be updated", 

212 ) 

213 if this_run.deployment_id is None: 

214 raise HTTPException( 

215 status_code=status.HTTP_400_BAD_REQUEST, 

216 detail="A deployment for the flow run could not be found", 

217 ) 

218 

219 deployment = await models.deployments.read_deployment( 

220 session=session, deployment_id=this_run.deployment_id 

221 ) 

222 if deployment is None: 

223 raise HTTPException( 

224 status_code=status.HTTP_400_BAD_REQUEST, 

225 detail="A deployment for the flow run could not be found", 

226 ) 

227 

228 await validate_job_variables_for_deployment_flow_run( 

229 session, deployment, flow_run 

230 ) 

231 

232 result = await models.flow_runs.update_flow_run( 

233 session=session, flow_run=flow_run, flow_run_id=flow_run_id 

234 ) 

235 if not result: 235 ↛ exitline 235 didn't return from function 'update_flow_run' because the condition on line 235 was always true

236 raise HTTPException(status.HTTP_404_NOT_FOUND, detail="Flow run not found") 

237 

238 

239@router.post("/count") 

240async def count_flow_runs( 

241 flows: Optional[schemas.filters.FlowFilter] = None, 

242 flow_runs: Optional[schemas.filters.FlowRunFilter] = None, 

243 task_runs: Optional[schemas.filters.TaskRunFilter] = None, 

244 deployments: Optional[schemas.filters.DeploymentFilter] = None, 

245 work_pools: Optional[schemas.filters.WorkPoolFilter] = None, 

246 work_pool_queues: Optional[schemas.filters.WorkQueueFilter] = None, 

247 db: PrefectDBInterface = Depends(provide_database_interface), 

248) -> int: 

249 """ 

250 Query for flow runs. 

251 """ 

252 async with db.session_context() as session: 

253 return await models.flow_runs.count_flow_runs( 

254 session=session, 

255 flow_filter=flows, 

256 flow_run_filter=flow_runs, 

257 task_run_filter=task_runs, 

258 deployment_filter=deployments, 

259 work_pool_filter=work_pools, 

260 work_queue_filter=work_pool_queues, 

261 ) 

262 

263 

264@router.post("/lateness") 

265async def average_flow_run_lateness( 

266 flows: Optional[schemas.filters.FlowFilter] = None, 

267 flow_runs: Optional[schemas.filters.FlowRunFilter] = None, 

268 task_runs: Optional[schemas.filters.TaskRunFilter] = None, 

269 deployments: Optional[schemas.filters.DeploymentFilter] = None, 

270 work_pools: Optional[schemas.filters.WorkPoolFilter] = None, 

271 work_pool_queues: Optional[schemas.filters.WorkQueueFilter] = None, 

272 db: PrefectDBInterface = Depends(provide_database_interface), 

273) -> Optional[float]: 

274 """ 

275 Query for average flow-run lateness in seconds. 

276 """ 

277 async with db.session_context() as session: 

278 if db.dialect.name == "sqlite": 278 ↛ 285line 278 didn't jump to line 285 because the condition on line 278 was never true

279 # Since we want an _average_ of the lateness we're unable to use 

280 # the existing FlowRun.expected_start_time_delta property as it 

281 # returns a timedelta and SQLite is unable to properly deal with it 

282 # and always returns 1970.0 as the average. This copies the same 

283 # logic but ensures that it returns the number of seconds instead 

284 # so it's compatible with SQLite. 

285 base_query = sa.case( 

286 ( 

287 db.FlowRun.start_time > db.FlowRun.expected_start_time, 

288 sa.func.strftime("%s", db.FlowRun.start_time) 

289 - sa.func.strftime("%s", db.FlowRun.expected_start_time), 

290 ), 

291 ( 

292 db.FlowRun.start_time.is_(None) 

293 & db.FlowRun.state_type.notin_(schemas.states.TERMINAL_STATES) 

294 & (db.FlowRun.expected_start_time < sa.func.datetime("now")), 

295 sa.func.strftime("%s", sa.func.datetime("now")) 

296 - sa.func.strftime("%s", db.FlowRun.expected_start_time), 

297 ), 

298 else_=0, 

299 ) 

300 else: 

301 base_query = db.FlowRun.estimated_start_time_delta 

302 

303 query = await models.flow_runs._apply_flow_run_filters( 

304 db, 

305 sa.select(sa.func.avg(base_query)), 

306 flow_filter=flows, 

307 flow_run_filter=flow_runs, 

308 task_run_filter=task_runs, 

309 deployment_filter=deployments, 

310 work_pool_filter=work_pools, 

311 work_queue_filter=work_pool_queues, 

312 ) 

313 result = await session.execute(query) 

314 

315 avg_lateness = result.scalar() 

316 

317 if avg_lateness is None: 

318 return None 

319 elif isinstance(avg_lateness, datetime.timedelta): 

320 return avg_lateness.total_seconds() 

321 else: 

322 return avg_lateness 

323 

324 

325@router.post("/history") 

326async def flow_run_history( 

327 history_start: DateTime = Body(..., description="The history's start time."), 

328 history_end: DateTime = Body(..., description="The history's end time."), 

329 # Workaround for the fact that FastAPI does not let us configure ser_json_timedelta 

330 # to represent timedeltas as floats in JSON. 

331 history_interval_seconds: float = Body( 

332 ..., 

333 description=( 

334 "The size of each history interval, in seconds. Must be at least 1 second." 

335 ), 

336 json_schema_extra={"format": "time-delta"}, 

337 ), 

338 flows: Optional[schemas.filters.FlowFilter] = None, 

339 flow_runs: Optional[schemas.filters.FlowRunFilter] = None, 

340 task_runs: Optional[schemas.filters.TaskRunFilter] = None, 

341 deployments: Optional[schemas.filters.DeploymentFilter] = None, 

342 work_pools: Optional[schemas.filters.WorkPoolFilter] = None, 

343 work_queues: Optional[schemas.filters.WorkQueueFilter] = None, 

344 db: PrefectDBInterface = Depends(provide_database_interface), 

345) -> List[schemas.responses.HistoryResponse]: 

346 """ 

347 Query for flow run history data across a given range and interval. 

348 """ 

349 history_interval = datetime.timedelta(seconds=history_interval_seconds) 

350 

351 if history_interval < datetime.timedelta(seconds=1): 351 ↛ 357line 351 didn't jump to line 357 because the condition on line 351 was always true

352 raise HTTPException( 

353 status.HTTP_422_UNPROCESSABLE_ENTITY, 

354 detail="History interval must not be less than 1 second.", 

355 ) 

356 

357 async with db.session_context() as session: 

358 return await run_history( 

359 session=session, 

360 run_type="flow_run", 

361 history_start=history_start, 

362 history_end=history_end, 

363 history_interval=history_interval, 

364 flows=flows, 

365 flow_runs=flow_runs, 

366 task_runs=task_runs, 

367 deployments=deployments, 

368 work_pools=work_pools, 

369 work_queues=work_queues, 

370 ) 

371 

372 

373@router.get("/{id:uuid}") 

374async def read_flow_run( 

375 flow_run_id: UUID = Path(..., description="The flow run id", alias="id"), 

376 db: PrefectDBInterface = Depends(provide_database_interface), 

377) -> schemas.responses.FlowRunResponse: 

378 """ 

379 Get a flow run by id. 

380 """ 

381 async with db.session_context() as session: 

382 flow_run = await models.flow_runs.read_flow_run( 

383 session=session, flow_run_id=flow_run_id 

384 ) 

385 if not flow_run: 

386 raise HTTPException(status.HTTP_404_NOT_FOUND, detail="Flow run not found") 

387 return schemas.responses.FlowRunResponse.model_validate( 

388 flow_run, from_attributes=True 

389 ) 

390 

391 

392@router.get("/{id:uuid}/graph", tags=["Flow Run Graph"]) 

393async def read_flow_run_graph_v1( 

394 flow_run_id: UUID = Path(..., description="The flow run id", alias="id"), 

395 db: PrefectDBInterface = Depends(provide_database_interface), 

396) -> List[DependencyResult]: 

397 """ 

398 Get a task run dependency map for a given flow run. 

399 """ 

400 async with db.session_context() as session: 

401 return await models.flow_runs.read_task_run_dependencies( 

402 session=session, flow_run_id=flow_run_id 

403 ) 

404 

405 

406@router.get("/{id:uuid}/graph-v2", tags=["Flow Run Graph"]) 

407async def read_flow_run_graph_v2( 

408 flow_run_id: UUID = Path(..., description="The flow run id", alias="id"), 

409 since: datetime.datetime = Query( 

410 default=jsonable_encoder(earliest_possible_datetime()), 

411 description="Only include runs that start or end after this time.", 

412 ), 

413 db: PrefectDBInterface = Depends(provide_database_interface), 

414) -> Graph: 

415 """ 

416 Get a graph of the tasks and subflow runs for the given flow run 

417 """ 

418 async with db.session_context() as session: 

419 try: 

420 return await read_flow_run_graph( 

421 session=session, 

422 flow_run_id=flow_run_id, 

423 since=since, 

424 ) 

425 except FlowRunGraphTooLarge as e: 

426 raise HTTPException( 

427 status_code=status.HTTP_400_BAD_REQUEST, 

428 detail=str(e), 

429 ) 

430 

431 

432@router.post("/{id:uuid}/resume") 

433async def resume_flow_run( 

434 response: Response, 

435 flow_run_id: UUID = Path(..., description="The flow run id", alias="id"), 

436 db: PrefectDBInterface = Depends(provide_database_interface), 

437 run_input: Optional[dict[str, Any]] = Body(default=None, embed=True), 

438 flow_policy: type[FlowRunOrchestrationPolicy] = Depends( 

439 orchestration_dependencies.provide_flow_policy 

440 ), 

441 task_policy: type[TaskRunOrchestrationPolicy] = Depends( 

442 orchestration_dependencies.provide_task_policy 

443 ), 

444 orchestration_parameters: Dict[str, Any] = Depends( 

445 orchestration_dependencies.provide_flow_orchestration_parameters 

446 ), 

447 api_version: str = Depends(dependencies.provide_request_api_version), 

448 client_version: Optional[str] = Depends(dependencies.get_prefect_client_version), 

449) -> OrchestrationResult: 

450 """ 

451 Resume a paused flow run. 

452 """ 

453 right_now = now("UTC") 

454 

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

456 flow_run = await models.flow_runs.read_flow_run(session, flow_run_id) 

457 state = flow_run.state 

458 

459 if state is None or state.type != schemas.states.StateType.PAUSED: 

460 result = OrchestrationResult( 

461 state=None, 

462 status=schemas.responses.SetStateStatus.ABORT, 

463 details=schemas.responses.StateAbortDetails( 

464 reason="Cannot resume a flow run that is not paused." 

465 ), 

466 ) 

467 return result 

468 

469 orchestration_parameters.update({"api-version": api_version}) 

470 

471 keyset = state.state_details.run_input_keyset 

472 

473 if keyset: 

474 run_input = run_input or {} 

475 

476 try: 

477 hydration_context = await schema_tools.HydrationContext.build( 

478 session=session, 

479 raise_on_error=True, 

480 render_jinja=True, 

481 render_workspace_variables=True, 

482 ) 

483 run_input = schema_tools.hydrate(run_input, hydration_context) or {} 

484 except schema_tools.HydrationError as exc: 

485 return OrchestrationResult( 

486 state=state, 

487 status=schemas.responses.SetStateStatus.REJECT, 

488 details=schemas.responses.StateAbortDetails( 

489 reason=f"Error hydrating run input: {exc}", 

490 ), 

491 ) 

492 

493 schema_json = await models.flow_run_input.read_flow_run_input( 

494 session=session, flow_run_id=flow_run.id, key=keyset["schema"] 

495 ) 

496 

497 if schema_json is None: 

498 return OrchestrationResult( 

499 state=state, 

500 status=schemas.responses.SetStateStatus.REJECT, 

501 details=schemas.responses.StateAbortDetails( 

502 reason="Run input schema not found." 

503 ), 

504 ) 

505 

506 try: 

507 schema = orjson.loads(schema_json.value) 

508 except orjson.JSONDecodeError: 

509 return OrchestrationResult( 

510 state=state, 

511 status=schemas.responses.SetStateStatus.REJECT, 

512 details=schemas.responses.StateAbortDetails( 

513 reason="Run input schema is not valid JSON." 

514 ), 

515 ) 

516 

517 try: 

518 schema_tools.validate(run_input, schema, raise_on_error=True) 

519 except schema_tools.ValidationError as exc: 

520 return OrchestrationResult( 

521 state=state, 

522 status=schemas.responses.SetStateStatus.REJECT, 

523 details=schemas.responses.StateAbortDetails( 

524 reason=f"Reason: {exc}" 

525 ), 

526 ) 

527 except schema_tools.CircularSchemaRefError: 

528 return OrchestrationResult( 

529 state=state, 

530 status=schemas.responses.SetStateStatus.REJECT, 

531 details=schemas.responses.StateAbortDetails( 

532 reason="Invalid schema: Unable to validate schema with circular references.", 

533 ), 

534 ) 

535 

536 if state.state_details.pause_reschedule: 

537 orchestration_result = await models.flow_runs.set_flow_run_state( 

538 session=session, 

539 flow_run_id=flow_run_id, 

540 state=schemas.states.Scheduled( 

541 name="Resuming", scheduled_time=now("UTC") 

542 ), 

543 flow_policy=flow_policy, 

544 orchestration_parameters=orchestration_parameters, 

545 client_version=client_version, 

546 ) 

547 else: 

548 orchestration_result = await models.flow_runs.set_flow_run_state( 

549 session=session, 

550 flow_run_id=flow_run_id, 

551 state=schemas.states.Running(), 

552 flow_policy=flow_policy, 

553 orchestration_parameters=orchestration_parameters, 

554 client_version=client_version, 

555 ) 

556 

557 if ( 

558 keyset 

559 and run_input 

560 and orchestration_result.status == schemas.responses.SetStateStatus.ACCEPT 

561 ): 

562 # The state change is accepted, go ahead and store the validated 

563 # run input. 

564 await models.flow_run_input.create_flow_run_input( 

565 session=session, 

566 flow_run_input=schemas.core.FlowRunInput( 

567 flow_run_id=flow_run_id, 

568 key=keyset["response"], 

569 value=orjson.dumps(run_input).decode("utf-8"), 

570 ), 

571 ) 

572 

573 # set the 201 if a new state was created 

574 if ( 

575 orchestration_result.state 

576 and orchestration_result.state.timestamp >= right_now 

577 ): 

578 response.status_code = status.HTTP_201_CREATED 

579 else: 

580 response.status_code = status.HTTP_200_OK 

581 

582 return orchestration_result 

583 

584 

585@router.post("/filter") 

586async def read_flow_runs( 

587 sort: schemas.sorting.FlowRunSort = Body(schemas.sorting.FlowRunSort.ID_DESC), 

588 limit: int = dependencies.LimitBody(), 

589 offset: int = Body(0, ge=0), 

590 flows: Optional[schemas.filters.FlowFilter] = None, 

591 flow_runs: Optional[schemas.filters.FlowRunFilter] = None, 

592 task_runs: Optional[schemas.filters.TaskRunFilter] = None, 

593 deployments: Optional[schemas.filters.DeploymentFilter] = None, 

594 work_pools: Optional[schemas.filters.WorkPoolFilter] = None, 

595 work_pool_queues: Optional[schemas.filters.WorkQueueFilter] = None, 

596 db: PrefectDBInterface = Depends(provide_database_interface), 

597) -> List[schemas.responses.FlowRunResponse]: 

598 """ 

599 Query for flow runs. 

600 """ 

601 async with db.session_context() as session: 

602 db_flow_runs = await models.flow_runs.read_flow_runs( 

603 session=session, 

604 flow_filter=flows, 

605 flow_run_filter=flow_runs, 

606 task_run_filter=task_runs, 

607 deployment_filter=deployments, 

608 work_pool_filter=work_pools, 

609 work_queue_filter=work_pool_queues, 

610 offset=offset, 

611 limit=limit, 

612 sort=sort, 

613 ) 

614 

615 # Instead of relying on fastapi.encoders.jsonable_encoder to convert the 

616 # response to JSON, we do so more efficiently ourselves. 

617 # In particular, the FastAPI encoder is very slow for large, nested objects. 

618 # See: https://github.com/tiangolo/fastapi/issues/1224 

619 encoded = [ 

620 schemas.responses.FlowRunResponse.model_validate( 

621 fr, from_attributes=True 

622 ).model_dump(mode="json") 

623 for fr in db_flow_runs 

624 ] 

625 return Response( 

626 content=orjson.dumps(encoded), 

627 media_type="application/json", 

628 ) 

629 

630 

631@router.delete("/{id:uuid}", status_code=status.HTTP_204_NO_CONTENT) 

632async def delete_flow_run( 

633 docket: dependencies.Docket, 

634 flow_run_id: UUID = Path(..., description="The flow run id", alias="id"), 

635 db: PrefectDBInterface = Depends(provide_database_interface), 

636) -> None: 

637 """ 

638 Delete a flow run by id. 

639 """ 

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

641 result = await models.flow_runs.delete_flow_run( 

642 session=session, flow_run_id=flow_run_id 

643 ) 

644 if not result: 

645 raise HTTPException( 

646 status_code=status.HTTP_404_NOT_FOUND, detail="Flow run not found" 

647 ) 

648 await docket.add( 

649 delete_flow_run_logs, 

650 key=f"delete_flow_run_logs:{flow_run_id}", 

651 )(flow_run_id=flow_run_id) 

652 

653 

654async def delete_flow_run_logs( 

655 *, 

656 db: PrefectDBInterface = DocketDepends(provide_database_interface), 

657 flow_run_id: UUID, 

658 retry: Retry = Retry(attempts=5, delay=datetime.timedelta(seconds=0.5)), 

659) -> None: 

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

661 await models.logs.delete_logs( 

662 session=session, 

663 log_filter=schemas.filters.LogFilter( 

664 flow_run_id=schemas.filters.LogFilterFlowRunId(any_=[flow_run_id]) 

665 ), 

666 ) 

667 

668 

669BULK_OPERATION_LIMIT = 50 

670 

671 

672@router.post("/bulk_delete") 

673async def bulk_delete_flow_runs( 

674 docket: dependencies.Docket, 

675 flow_runs: Optional[schemas.filters.FlowRunFilter] = Body( 

676 None, description="Filter criteria for flow runs to delete" 

677 ), 

678 limit: int = Body( 

679 BULK_OPERATION_LIMIT, 

680 ge=1, 

681 le=BULK_OPERATION_LIMIT, 

682 description=f"Maximum number of flow runs to delete. Defaults to {BULK_OPERATION_LIMIT}.", 

683 ), 

684 db: PrefectDBInterface = Depends(provide_database_interface), 

685) -> FlowRunBulkDeleteResponse: 

686 """ 

687 Bulk delete flow runs matching the specified filter criteria. 

688 

689 Returns the IDs of flow runs that were deleted. 

690 """ 

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

692 # Query matching flow runs 

693 db_flow_runs = await models.flow_runs.read_flow_runs( 

694 session=session, 

695 flow_run_filter=flow_runs, 

696 limit=limit, 

697 ) 

698 

699 if not db_flow_runs: 

700 return FlowRunBulkDeleteResponse(deleted=[]) 

701 

702 flow_run_ids = [fr.id for fr in db_flow_runs] 

703 

704 # Delete flow runs 

705 deleted_ids = await models.flow_runs.delete_flow_runs( 

706 session=session, 

707 flow_run_ids=flow_run_ids, 

708 ) 

709 

710 # Queue log cleanup for each deleted flow run 

711 for flow_run_id in deleted_ids: 

712 await docket.add( 

713 delete_flow_run_logs, 

714 key=f"delete_flow_run_logs:{flow_run_id}", 

715 )(flow_run_id=flow_run_id) 

716 

717 return FlowRunBulkDeleteResponse(deleted=deleted_ids) 

718 

719 

720@router.post("/bulk_set_state") 

721async def bulk_set_flow_run_state( 

722 request: Request, 

723 flow_runs: Optional[schemas.filters.FlowRunFilter] = Body( 

724 None, description="Filter criteria for flow runs to update" 

725 ), 

726 state: schemas.actions.StateCreate = Body(..., description="The state to set"), 

727 force: bool = Body( 

728 False, 

729 description=( 

730 "If false, orchestration rules will be applied that may alter or prevent" 

731 " the state transition. If True, orchestration rules are not applied." 

732 ), 

733 ), 

734 limit: int = Body( 

735 BULK_OPERATION_LIMIT, 

736 ge=1, 

737 le=BULK_OPERATION_LIMIT, 

738 description=f"Maximum number of flow runs to update. Defaults to {BULK_OPERATION_LIMIT}.", 

739 ), 

740 db: PrefectDBInterface = Depends(provide_database_interface), 

741 flow_policy: type[FlowRunOrchestrationPolicy] = Depends( 

742 orchestration_dependencies.provide_flow_policy 

743 ), 

744 orchestration_parameters: Dict[str, Any] = Depends( 

745 orchestration_dependencies.provide_flow_orchestration_parameters 

746 ), 

747 api_version: str = Depends(dependencies.provide_request_api_version), 

748 client_version: Optional[str] = Depends(dependencies.get_prefect_client_version), 

749) -> FlowRunBulkSetStateResponse: 

750 """ 

751 Bulk set state for flow runs matching the specified filter criteria. 

752 

753 Returns the orchestration results for each flow run. 

754 """ 

755 orchestration_parameters.update({"api-version": api_version}) 

756 

757 async with db.session_context() as session: 

758 # Query matching flow runs 

759 db_flow_runs = await models.flow_runs.read_flow_runs( 

760 session=session, 

761 flow_run_filter=flow_runs, 

762 limit=limit, 

763 ) 

764 

765 if not db_flow_runs: 

766 return FlowRunBulkSetStateResponse(results=[]) 

767 

768 results: List[FlowRunOrchestrationResult] = [] 

769 

770 # Process flow runs sequentially to avoid session conflicts 

771 for flow_run in db_flow_runs: 

772 state_to_schedule: schemas.states.State | None = None 

773 async with db.session_context( 

774 begin_transaction=True, with_for_update=True 

775 ) as session: 

776 try: 

777 orchestration_result = await models.flow_runs.set_flow_run_state( 

778 session=session, 

779 flow_run_id=flow_run.id, 

780 state=schemas.states.State.model_validate(state), 

781 force=force, 

782 flow_policy=flow_policy, 

783 orchestration_parameters=orchestration_parameters, 

784 client_version=client_version, 

785 ) 

786 results.append( 

787 FlowRunOrchestrationResult( 

788 flow_run_id=flow_run.id, 

789 status=orchestration_result.status, 

790 state=orchestration_result.state, 

791 details=orchestration_result.details, 

792 ) 

793 ) 

794 if ( 

795 orchestration_result.status 

796 == schemas.responses.SetStateStatus.ACCEPT 

797 ): 

798 state_to_schedule = orchestration_result.state 

799 except Exception as e: 

800 results.append( 

801 FlowRunOrchestrationResult( 

802 flow_run_id=flow_run.id, 

803 status=schemas.responses.SetStateStatus.ABORT, 

804 state=None, 

805 details=schemas.responses.StateAbortDetails(reason=str(e)), 

806 ) 

807 ) 

808 continue 

809 

810 if state_to_schedule is not None: 

811 await _maybe_schedule_cancelling_timeout_check_for_state( 

812 request=request, 

813 flow_run_id=flow_run.id, 

814 state=state_to_schedule, 

815 ) 

816 

817 return FlowRunBulkSetStateResponse(results=results) 

818 

819 

820@router.post("/{id:uuid}/set_state") 

821async def set_flow_run_state( 

822 response: Response, 

823 request: Request, 

824 flow_run_id: UUID = Path(..., description="The flow run id", alias="id"), 

825 state: schemas.actions.StateCreate = Body(..., description="The intended state."), 

826 force: bool = Body( 

827 False, 

828 description=( 

829 "If false, orchestration rules will be applied that may alter or prevent" 

830 " the state transition. If True, orchestration rules are not applied." 

831 ), 

832 ), 

833 db: PrefectDBInterface = Depends(provide_database_interface), 

834 flow_policy: type[FlowRunOrchestrationPolicy] = Depends( 

835 orchestration_dependencies.provide_flow_policy 

836 ), 

837 orchestration_parameters: Dict[str, Any] = Depends( 

838 orchestration_dependencies.provide_flow_orchestration_parameters 

839 ), 

840 api_version: str = Depends(dependencies.provide_request_api_version), 

841 client_version: Optional[str] = Depends(dependencies.get_prefect_client_version), 

842) -> OrchestrationResult: 

843 """Set a flow run state, invoking any orchestration rules.""" 

844 

845 # pass the request version to the orchestration engine to support compatibility code 

846 orchestration_parameters.update({"api-version": api_version}) 

847 

848 right_now = now("UTC") 

849 

850 # create the state 

851 async with db.session_context( 

852 begin_transaction=True, with_for_update=True 

853 ) as session: 

854 orchestration_result = await models.flow_runs.set_flow_run_state( 

855 session=session, 

856 flow_run_id=flow_run_id, 

857 # convert to a full State object 

858 state=schemas.states.State.model_validate(state), 

859 force=force, 

860 flow_policy=flow_policy, 

861 orchestration_parameters=orchestration_parameters, 

862 client_version=client_version, 

863 ) 

864 

865 if orchestration_result.status == schemas.responses.SetStateStatus.ACCEPT: 

866 await _maybe_schedule_cancelling_timeout_check_for_state( 

867 request=request, 

868 flow_run_id=flow_run_id, 

869 state=orchestration_result.state, 

870 ) 

871 

872 # set the 201 if a new state was created 

873 if orchestration_result.state and orchestration_result.state.timestamp >= right_now: 

874 response.status_code = status.HTTP_201_CREATED 

875 else: 

876 response.status_code = status.HTTP_200_OK 

877 

878 return orchestration_result 

879 

880 

881@router.post("/{id:uuid}/input", status_code=status.HTTP_201_CREATED) 

882async def create_flow_run_input( 

883 flow_run_id: UUID = Path(..., description="The flow run id", alias="id"), 

884 key: str = Body(..., description="The input key"), 

885 value: bytes = Body(..., description="The value of the input"), 

886 sender: Optional[str] = Body(None, description="The sender of the input"), 

887 db: PrefectDBInterface = Depends(provide_database_interface), 

888) -> None: 

889 """ 

890 Create a key/value input for a flow run. 

891 """ 

892 async with db.session_context() as session: 

893 try: 

894 await models.flow_run_input.create_flow_run_input( 

895 session=session, 

896 flow_run_input=schemas.core.FlowRunInput( 

897 flow_run_id=flow_run_id, 

898 key=key, 

899 sender=sender, 

900 value=value.decode(), 

901 ), 

902 ) 

903 await session.commit() 

904 

905 except IntegrityError as exc: 

906 if "unique constraint" in str(exc).lower(): 

907 raise HTTPException( 

908 status_code=status.HTTP_409_CONFLICT, 

909 detail="A flow run input with this key already exists.", 

910 ) 

911 else: 

912 raise HTTPException( 

913 status_code=status.HTTP_404_NOT_FOUND, detail="Flow run not found" 

914 ) 

915 

916 

917@router.post("/{id:uuid}/input/filter") 

918async def filter_flow_run_input( 

919 flow_run_id: UUID = Path(..., description="The flow run id", alias="id"), 

920 prefix: str = Body(..., description="The input key prefix", embed=True), 

921 limit: int = Body( 

922 1, description="The maximum number of results to return", embed=True 

923 ), 

924 exclude_keys: List[str] = Body( 

925 [], description="Exclude inputs with these keys", embed=True 

926 ), 

927 db: PrefectDBInterface = Depends(provide_database_interface), 

928) -> List[schemas.core.FlowRunInput]: 

929 """ 

930 Filter flow run inputs by key prefix 

931 """ 

932 async with db.session_context() as session: 

933 return await models.flow_run_input.filter_flow_run_input( 

934 session=session, 

935 flow_run_id=flow_run_id, 

936 prefix=prefix, 

937 limit=limit, 

938 exclude_keys=exclude_keys, 

939 ) 

940 

941 

942@router.get("/{id:uuid}/input/{key}") 

943async def read_flow_run_input( 

944 flow_run_id: UUID = Path(..., description="The flow run id", alias="id"), 

945 key: str = Path(..., description="The input key", alias="key"), 

946 db: PrefectDBInterface = Depends(provide_database_interface), 

947) -> PlainTextResponse: 

948 """ 

949 Create a value from a flow run input 

950 """ 

951 

952 async with db.session_context() as session: 

953 flow_run_input = await models.flow_run_input.read_flow_run_input( 

954 session=session, flow_run_id=flow_run_id, key=key 

955 ) 

956 

957 if flow_run_input: 957 ↛ 958line 957 didn't jump to line 958 because the condition on line 957 was never true

958 return PlainTextResponse(flow_run_input.value) 

959 else: 

960 raise HTTPException( 

961 status_code=status.HTTP_404_NOT_FOUND, detail="Flow run input not found" 

962 ) 

963 

964 

965@router.delete("/{id:uuid}/input/{key}", status_code=status.HTTP_204_NO_CONTENT) 

966async def delete_flow_run_input( 

967 flow_run_id: UUID = Path(..., description="The flow run id", alias="id"), 

968 key: str = Path(..., description="The input key", alias="key"), 

969 db: PrefectDBInterface = Depends(provide_database_interface), 

970) -> None: 

971 """ 

972 Delete a flow run input 

973 """ 

974 

975 async with db.session_context() as session: 

976 deleted = await models.flow_run_input.delete_flow_run_input( 

977 session=session, flow_run_id=flow_run_id, key=key 

978 ) 

979 await session.commit() 

980 

981 if not deleted: 

982 raise HTTPException( 

983 status_code=status.HTTP_404_NOT_FOUND, detail="Flow run input not found" 

984 ) 

985 

986 

987@router.post("/paginate") 

988async def paginate_flow_runs( 

989 sort: schemas.sorting.FlowRunSort = Body(schemas.sorting.FlowRunSort.ID_DESC), 

990 limit: int = dependencies.LimitBody(), 

991 page: int = Body(1, ge=1), 

992 flows: Optional[schemas.filters.FlowFilter] = None, 

993 flow_runs: Optional[schemas.filters.FlowRunFilter] = None, 

994 task_runs: Optional[schemas.filters.TaskRunFilter] = None, 

995 deployments: Optional[schemas.filters.DeploymentFilter] = None, 

996 work_pools: Optional[schemas.filters.WorkPoolFilter] = None, 

997 work_pool_queues: Optional[schemas.filters.WorkQueueFilter] = None, 

998 db: PrefectDBInterface = Depends(provide_database_interface), 

999) -> FlowRunPaginationResponse: 

1000 """ 

1001 Pagination query for flow runs. 

1002 """ 

1003 offset = (page - 1) * limit 

1004 

1005 async def get_runs(): 

1006 async with db.session_context() as session: 

1007 return await models.flow_runs.read_flow_runs( 

1008 session=session, 

1009 flow_filter=flows, 

1010 flow_run_filter=flow_runs, 

1011 task_run_filter=task_runs, 

1012 deployment_filter=deployments, 

1013 work_pool_filter=work_pools, 

1014 work_queue_filter=work_pool_queues, 

1015 offset=offset, 

1016 limit=limit, 

1017 sort=sort, 

1018 ) 

1019 

1020 async def get_count(): 

1021 async with db.session_context() as session: 

1022 return await models.flow_runs.count_flow_runs( 

1023 session=session, 

1024 flow_filter=flows, 

1025 flow_run_filter=flow_runs, 

1026 task_run_filter=task_runs, 

1027 deployment_filter=deployments, 

1028 work_pool_filter=work_pools, 

1029 work_queue_filter=work_pool_queues, 

1030 ) 

1031 

1032 runs, count = await asyncio.gather(get_runs(), get_count()) 

1033 

1034 # Instead of relying on fastapi.encoders.jsonable_encoder to convert the 

1035 # response to JSON, we do so more efficiently ourselves. 

1036 # In particular, the FastAPI encoder is very slow for large, nested objects. 

1037 # See: https://github.com/tiangolo/fastapi/issues/1224 

1038 results = [ 

1039 schemas.responses.FlowRunResponse.model_validate(run, from_attributes=True) 

1040 for run in runs 

1041 ] 

1042 

1043 response = FlowRunPaginationResponse( 

1044 results=results, 

1045 count=count, 

1046 limit=limit, 

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

1048 page=page, 

1049 ).model_dump(mode="json") 

1050 

1051 return Response( 

1052 content=orjson.dumps(response), 

1053 media_type="application/json", 

1054 ) 

1055 

1056 

1057FLOW_RUN_LOGS_DOWNLOAD_PAGE_LIMIT = 1000 

1058 

1059 

1060@router.get("/{id:uuid}/logs/download") 

1061async def download_logs( 

1062 flow_run_id: UUID = Path(..., description="The flow run id", alias="id"), 

1063 db: PrefectDBInterface = Depends(provide_database_interface), 

1064) -> StreamingResponse: 

1065 """ 

1066 Download all flow run logs as a CSV file, collecting all logs until there are no more logs to retrieve. 

1067 """ 

1068 async with db.session_context() as session: 

1069 flow_run = await models.flow_runs.read_flow_run( 

1070 session=session, flow_run_id=flow_run_id 

1071 ) 

1072 

1073 if not flow_run: 

1074 raise HTTPException(status.HTTP_404_NOT_FOUND, detail="Flow run not found") 

1075 

1076 filename = quote(f"{flow_run.name}-logs.csv", safe="") 

1077 

1078 async def generate(): 

1079 

1080 data = io.StringIO() 

1081 csv_writer = csv.writer(data) 

1082 

1083 csv_writer.writerow( 

1084 ["timestamp", "level", "flow_run_id", "task_run_id", "message"] 

1085 ) 

1086 data.seek(0) 

1087 yield data.read() 

1088 data.seek(0) 

1089 data.truncate(0) 

1090 

1091 offset = 0 

1092 limit = FLOW_RUN_LOGS_DOWNLOAD_PAGE_LIMIT 

1093 

1094 async with db.session_context() as session: 

1095 while True: 

1096 results = await models.logs.read_logs( 

1097 session=session, 

1098 log_filter=schemas.filters.LogFilter( 

1099 flow_run_id={"any_": [flow_run_id]} 

1100 ), 

1101 offset=offset, 

1102 limit=limit, 

1103 sort=schemas.sorting.LogSort.TIMESTAMP_ASC, 

1104 ) 

1105 

1106 if not results: 

1107 break 

1108 

1109 offset += limit 

1110 

1111 for log in results: 

1112 csv_writer.writerow( 

1113 [ 

1114 log.timestamp, 

1115 log.level, 

1116 log.flow_run_id, 

1117 log.task_run_id, 

1118 log.message, 

1119 ] 

1120 ) 

1121 data.seek(0) 

1122 yield data.read() 

1123 data.seek(0) 

1124 data.truncate(0) 

1125 

1126 return StreamingResponse( 

1127 generate(), 

1128 media_type="text/csv; charset=utf-8", 

1129 headers={ 

1130 "Content-Disposition": ( 

1131 f"attachment; filename=\"flow-run-logs.csv\"; filename*=UTF-8''{filename}" 

1132 ) 

1133 }, 

1134 ) 

1135 

1136 

1137@router.patch("/{id:uuid}/labels", status_code=status.HTTP_204_NO_CONTENT) 

1138async def update_flow_run_labels( 

1139 flow_run_id: UUID = Path(..., description="The flow run id", alias="id"), 

1140 labels: Dict[str, Any] = Body(..., description="The labels to update"), 

1141 db: PrefectDBInterface = Depends(provide_database_interface), 

1142) -> None: 

1143 """ 

1144 Update the labels of a flow run. 

1145 """ 

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

1147 await models.flow_runs.update_flow_run_labels( 

1148 session=session, flow_run_id=flow_run_id, labels=labels 

1149 )