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
« 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"""
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
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
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
65if TYPE_CHECKING: 65 ↛ 66line 65 didn't jump to line 66 because the condition on line 65 was never true
66 import logging
68logger: "logging.Logger" = get_logger("server.api")
70router: PrefectRouter = PrefectRouter(prefix="/flow_runs", tags=["Flow Runs"])
73def _get_request_docket(request: Request) -> Docket | None:
74 return getattr(request.app.state, "docket", None)
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
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 )
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.
121 If no state is provided, the flow run will be created in a PENDING state.
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 )
130 # pass the request version to the orchestration engine to support compatibility code
131 orchestration_parameters.update({"api-version": api_version})
133 if not flow_run_object.state:
134 flow_run_object.state = schemas.states.Pending()
136 right_now = now("UTC")
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
156 flow_run_object.work_queue_id = work_queue_id
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 )
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
182 return flow_run_response
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 )
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 )
228 await validate_job_variables_for_deployment_flow_run(
229 session, deployment, flow_run
230 )
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")
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 )
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
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)
315 avg_lateness = result.scalar()
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
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)
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 )
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 )
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 )
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 )
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 )
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")
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
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
469 orchestration_parameters.update({"api-version": api_version})
471 keyset = state.state_details.run_input_keyset
473 if keyset:
474 run_input = run_input or {}
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 )
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 )
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 )
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 )
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 )
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 )
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 )
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
582 return orchestration_result
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 )
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 )
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)
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 )
669BULK_OPERATION_LIMIT = 50
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.
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 )
699 if not db_flow_runs:
700 return FlowRunBulkDeleteResponse(deleted=[])
702 flow_run_ids = [fr.id for fr in db_flow_runs]
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 )
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)
717 return FlowRunBulkDeleteResponse(deleted=deleted_ids)
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.
753 Returns the orchestration results for each flow run.
754 """
755 orchestration_parameters.update({"api-version": api_version})
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 )
765 if not db_flow_runs:
766 return FlowRunBulkSetStateResponse(results=[])
768 results: List[FlowRunOrchestrationResult] = []
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
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 )
817 return FlowRunBulkSetStateResponse(results=results)
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."""
845 # pass the request version to the orchestration engine to support compatibility code
846 orchestration_parameters.update({"api-version": api_version})
848 right_now = now("UTC")
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 )
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 )
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
878 return orchestration_result
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()
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 )
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 )
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 """
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 )
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 )
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 """
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()
981 if not deleted:
982 raise HTTPException(
983 status_code=status.HTTP_404_NOT_FOUND, detail="Flow run input not found"
984 )
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
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 )
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 )
1032 runs, count = await asyncio.gather(get_runs(), get_count())
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 ]
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")
1051 return Response(
1052 content=orjson.dumps(response),
1053 media_type="application/json",
1054 )
1057FLOW_RUN_LOGS_DOWNLOAD_PAGE_LIMIT = 1000
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 )
1073 if not flow_run:
1074 raise HTTPException(status.HTTP_404_NOT_FOUND, detail="Flow run not found")
1076 filename = quote(f"{flow_run.name}-logs.csv", safe="")
1078 async def generate():
1080 data = io.StringIO()
1081 csv_writer = csv.writer(data)
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)
1091 offset = 0
1092 limit = FLOW_RUN_LOGS_DOWNLOAD_PAGE_LIMIT
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 )
1106 if not results:
1107 break
1109 offset += limit
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)
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 )
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 )