Coverage for /usr/local/lib/python3.12/site-packages/prefect/server/api/deployments.py: 35%
387 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 Deployment objects.
3"""
5import datetime
6import logging
7from typing import List, Optional
8from uuid import UUID
10import jsonschema.exceptions
11import sqlalchemy as sa
12from fastapi import Body, Depends, HTTPException, Path, Response
14import prefect.server.api.dependencies as dependencies
15import prefect.server.models as models
16import prefect.server.schemas as schemas
17from prefect._internal.compatibility.starlette import status
18from prefect.logging import get_logger
19from prefect.server.api.validation import (
20 validate_job_variables_for_deployment,
21 validate_job_variables_for_deployment_flow_run,
22)
23from prefect.server.api.workers import WorkerLookups
24from prefect.server.database import PrefectDBInterface, provide_database_interface
25from prefect.server.exceptions import MissingVariableError, ObjectNotFoundError
26from prefect.server.models.deployments import mark_deployments_ready
27from prefect.server.models.workers import DEFAULT_AGENT_WORK_POOL_NAME
28from prefect.server.schemas.responses import (
29 DeploymentBulkDeleteResponse,
30 DeploymentPaginationResponse,
31 FlowRunBulkCreateResponse,
32 FlowRunCreateResult,
33)
34from prefect.server.utilities.server import PrefectRouter
35from prefect.types import DateTime
36from prefect.types._datetime import now
37from prefect.utilities.schema_tools.hydration import (
38 HydrationContext,
39 HydrationError,
40 hydrate,
41)
42from prefect.utilities.schema_tools.validation import (
43 CircularSchemaRefError,
44 ValidationError,
45 validate,
46)
48logger: logging.Logger = get_logger(__name__)
50router: PrefectRouter = PrefectRouter(prefix="/deployments", tags=["Deployments"])
53def _multiple_schedules_error(deployment_id) -> HTTPException:
54 return HTTPException(
55 status.HTTP_422_UNPROCESSABLE_ENTITY,
56 detail=(
57 "Error updating deployment: "
58 f"Deployment {deployment_id!r} has multiple schedules. "
59 "Please use the UI or update your client to adjust this "
60 "deployment's schedules.",
61 ),
62 )
65@router.post("/")
66async def create_deployment(
67 deployment: schemas.actions.DeploymentCreate,
68 response: Response,
69 worker_lookups: WorkerLookups = Depends(WorkerLookups),
70 created_by: Optional[schemas.core.CreatedBy] = Depends(dependencies.get_created_by),
71 updated_by: Optional[schemas.core.UpdatedBy] = Depends(dependencies.get_updated_by),
72 db: PrefectDBInterface = Depends(provide_database_interface),
73) -> schemas.responses.DeploymentResponse:
74 """
75 Creates a new deployment from the provided schema. If a deployment with
76 the same name and flow_id already exists, the deployment is updated.
78 If the deployment has an active schedule, flow runs will be scheduled.
79 When upserting, any scheduled runs from the existing deployment will be deleted.
81 For more information, see https://docs.prefect.io/v3/concepts/deployments.
82 """
84 data = deployment.model_dump(exclude_unset=True)
85 data["created_by"] = created_by.model_dump() if created_by else None
86 data["updated_by"] = updated_by.model_dump() if created_by else None
88 async with db.session_context(begin_transaction=True) as session:
89 if (
90 deployment.work_pool_name
91 and deployment.work_pool_name != DEFAULT_AGENT_WORK_POOL_NAME
92 ):
93 # Make sure that deployment is valid before beginning creation process
94 work_pool = await models.workers.read_work_pool_by_name(
95 session=session, work_pool_name=deployment.work_pool_name
96 )
97 if work_pool is None:
98 raise HTTPException(
99 status_code=status.HTTP_404_NOT_FOUND,
100 detail=f'Work pool "{deployment.work_pool_name}" not found.',
101 )
103 await validate_job_variables_for_deployment(
104 session,
105 work_pool,
106 deployment,
107 )
109 # hydrate the input model into a full model
110 deployment_dict: dict = deployment.model_dump(
111 exclude={"work_pool_name"},
112 exclude_unset=True,
113 )
115 requested_concurrency_limit = deployment_dict.pop(
116 "global_concurrency_limit_id", "unset"
117 )
118 if requested_concurrency_limit != "unset":
119 if requested_concurrency_limit:
120 concurrency_limit = (
121 await models.concurrency_limits_v2.read_concurrency_limit(
122 session=session,
123 concurrency_limit_id=requested_concurrency_limit,
124 )
125 )
127 if not concurrency_limit:
128 raise HTTPException(
129 status_code=status.HTTP_404_NOT_FOUND,
130 detail="Concurrency limit not found",
131 )
133 deployment_dict["concurrency_limit_id"] = requested_concurrency_limit
135 if deployment.work_pool_name and deployment.work_queue_name: 135 ↛ 138line 135 didn't jump to line 138 because the condition on line 135 was never true
136 # If a specific pool name/queue name combination was provided, get the
137 # ID for that work pool queue.
138 deployment_dict[
139 "work_queue_id"
140 ] = await worker_lookups._get_work_queue_id_from_name(
141 session=session,
142 work_pool_name=deployment.work_pool_name,
143 work_queue_name=deployment.work_queue_name,
144 create_queue_if_not_found=True,
145 )
146 elif deployment.work_pool_name: 146 ↛ 149line 146 didn't jump to line 149 because the condition on line 146 was never true
147 # If just a pool name was provided, get the ID for its default
148 # work pool queue.
149 deployment_dict[
150 "work_queue_id"
151 ] = await worker_lookups._get_default_work_queue_id_from_work_pool_name(
152 session=session,
153 work_pool_name=deployment.work_pool_name,
154 )
155 elif deployment.work_queue_name:
156 # If just a queue name was provided, ensure that the queue exists and
157 # get its ID.
158 work_queue = await models.work_queues.ensure_work_queue_exists(
159 session=session, name=deployment.work_queue_name
160 )
161 deployment_dict["work_queue_id"] = work_queue.id
163 deployment = schemas.core.Deployment(**deployment_dict)
164 # check to see if relevant blocks exist, allowing us throw a useful error message
165 # for debugging
166 if deployment.infrastructure_document_id is not None:
167 infrastructure_block = (
168 await models.block_documents.read_block_document_by_id(
169 session=session,
170 block_document_id=deployment.infrastructure_document_id,
171 )
172 )
173 if not infrastructure_block: 173 ↛ 184line 173 didn't jump to line 184 because the condition on line 173 was always true
174 raise HTTPException(
175 status_code=status.HTTP_409_CONFLICT,
176 detail=(
177 "Error creating deployment. Could not find infrastructure"
178 f" block with id: {deployment.infrastructure_document_id}. This"
179 " usually occurs when applying a deployment specification that"
180 " was built against a different Prefect database / workspace."
181 ),
182 )
184 if deployment.storage_document_id is not None:
185 storage_block = await models.block_documents.read_block_document_by_id(
186 session=session,
187 block_document_id=deployment.storage_document_id,
188 )
189 if not storage_block: 189 ↛ anywhereline 189 didn't jump anywhere: it always raised an exception.
190 raise HTTPException(
191 status_code=status.HTTP_409_CONFLICT,
192 detail=(
193 "Error creating deployment. Could not find storage block with"
194 f" id: {deployment.storage_document_id}. This usually occurs"
195 " when applying a deployment specification that was built"
196 " against a different Prefect database / workspace."
197 ),
198 )
200 right_now = now("UTC")
201 model = await models.deployments.create_deployment(
202 session=session, deployment=deployment
203 )
205 if model.created >= right_now: 205 ↛ anywhereline 205 didn't jump anywhere: it always raised an exception.
206 response.status_code = status.HTTP_201_CREATED
208 return schemas.responses.DeploymentResponse.model_validate(
209 model, from_attributes=True
210 )
213@router.patch("/{id:uuid}", status_code=status.HTTP_204_NO_CONTENT)
214async def update_deployment(
215 deployment: schemas.actions.DeploymentUpdate,
216 deployment_id: UUID = Path(..., description="The deployment id", alias="id"),
217 db: PrefectDBInterface = Depends(provide_database_interface),
218) -> None:
219 async with db.session_context(begin_transaction=True) as session:
220 existing_deployment = await models.deployments.read_deployment(
221 session=session, deployment_id=deployment_id
222 )
223 if not existing_deployment:
224 raise HTTPException(
225 status.HTTP_404_NOT_FOUND, detail="Deployment not found."
226 )
228 # Checking how we should handle schedule updates
229 # If not all existing schedules have slugs then we'll fall back to the existing logic where are schedules are recreated to match the request.
230 # If the existing schedules have slugs, but not all provided schedules have slugs, then we'll return a 422 to avoid accidentally blowing away schedules.
231 # Otherwise, we'll use the existing slugs and the provided slugs to make targeted updates to the deployment's schedules.
232 schedules_to_patch: list[schemas.actions.DeploymentScheduleUpdate] = []
233 schedules_to_create: list[schemas.actions.DeploymentScheduleUpdate] = []
234 all_provided_have_slugs = all(
235 schedule.slug is not None for schedule in deployment.schedules or []
236 )
237 all_existing_have_slugs = existing_deployment.schedules and all(
238 schedule.slug is not None for schedule in existing_deployment.schedules
239 )
240 if all_provided_have_slugs and all_existing_have_slugs:
241 current_slugs = [
242 schedule.slug for schedule in existing_deployment.schedules
243 ]
245 # Check for duplicate replaces targets
246 replaces_targets: dict[str, str] = {} # old_slug -> new_slug
247 for schedule in deployment.schedules or []:
248 if schedule.replaces:
249 if schedule.replaces in replaces_targets:
250 raise HTTPException(
251 status_code=status.HTTP_422_UNPROCESSABLE_ENTITY,
252 detail=f"Multiple schedules have 'replaces' targeting the same slug: {schedule.replaces}",
253 )
254 replaces_targets[schedule.replaces] = schedule.slug or ""
256 # Check for slug collisions: if a schedule's new slug already
257 # exists and that existing slug is not being replaced, it's a
258 # collision
259 slugs_being_replaced = set(replaces_targets.keys())
260 for schedule in deployment.schedules:
261 if (
262 schedule.slug
263 and schedule.slug in current_slugs
264 and schedule.slug not in slugs_being_replaced
265 and schedule.replaces
266 ):
267 raise HTTPException(
268 status_code=status.HTTP_422_UNPROCESSABLE_ENTITY,
269 detail=f"Cannot rename schedule from '{schedule.replaces}' to '{schedule.slug}': "
270 f"a schedule with slug '{schedule.slug}' already exists.",
271 )
273 for schedule in deployment.schedules:
274 # Check if this schedule replaces an existing one
275 target_slug = schedule.replaces if schedule.replaces else schedule.slug
277 if target_slug in current_slugs: 277 ↛ 279line 277 didn't jump to line 279 because the condition on line 277 was always true
278 schedules_to_patch.append(schedule)
279 elif schedule.replaces:
280 # replaces points to a non-existent slug - warn and create new
281 logger.warning(
282 f"Schedule with slug '{schedule.slug}' has 'replaces: {schedule.replaces}' "
283 f"but no schedule with slug '{schedule.replaces}' exists. Creating new schedule."
284 )
285 if schedule.schedule:
286 schedules_to_create.append(schedule)
287 else:
288 raise HTTPException(
289 status_code=status.HTTP_422_UNPROCESSABLE_ENTITY,
290 detail="Unable to create new deployment schedules without a schedule configuration.",
291 )
292 elif schedule.schedule:
293 schedules_to_create.append(schedule)
294 else:
295 raise HTTPException(
296 status_code=status.HTTP_422_UNPROCESSABLE_ENTITY,
297 detail="Unable to create new deployment schedules without a schedule configuration.",
298 )
299 # Clear schedules to handle their update/creation separately
300 deployment.schedules = None
301 elif not all_provided_have_slugs and all_existing_have_slugs:
302 raise HTTPException(
303 status_code=status.HTTP_422_UNPROCESSABLE_ENTITY,
304 detail="Please provide a slug for each schedule in your request to ensure schedules are updated correctly.",
305 )
307 if deployment.work_pool_name:
308 # Make sure that deployment is valid before beginning creation process
309 work_pool = await models.workers.read_work_pool_by_name(
310 session=session, work_pool_name=deployment.work_pool_name
311 )
312 try:
313 deployment.check_valid_configuration(work_pool.base_job_template)
314 except (
315 MissingVariableError,
316 jsonschema.exceptions.ValidationError,
317 ValidationError,
318 ) as exc:
319 raise HTTPException(
320 status_code=status.HTTP_409_CONFLICT,
321 detail=f"Error creating deployment: {exc!r}",
322 )
324 if deployment.parameters is not None:
325 try:
326 dehydrated_params = deployment.parameters
327 ctx = await HydrationContext.build(
328 session=session,
329 raise_on_error=True,
330 render_jinja=True,
331 render_workspace_variables=True,
332 )
333 parameters = hydrate(dehydrated_params, ctx)
334 deployment.parameters = parameters
335 except HydrationError as exc:
336 raise HTTPException(
337 status.HTTP_400_BAD_REQUEST,
338 detail=f"Error hydrating deployment parameters: {exc}",
339 )
340 else:
341 parameters = existing_deployment.parameters
343 enforce_parameter_schema = (
344 deployment.enforce_parameter_schema
345 if deployment.enforce_parameter_schema is not None
346 else existing_deployment.enforce_parameter_schema
347 )
348 if enforce_parameter_schema:
349 # ensure that the new parameters conform to the proposed schema
350 if deployment.parameter_openapi_schema:
351 openapi_schema = deployment.parameter_openapi_schema
352 else:
353 openapi_schema = existing_deployment.parameter_openapi_schema
355 if not isinstance(openapi_schema, dict):
356 raise HTTPException(
357 status.HTTP_409_CONFLICT,
358 detail=(
359 "Error updating deployment: Cannot update parameters because"
360 " parameter schema enforcement is enabled and the deployment"
361 " does not have a valid parameter schema."
362 ),
363 )
364 try:
365 validate(
366 parameters,
367 openapi_schema,
368 raise_on_error=True,
369 ignore_required=True,
370 )
371 except ValidationError as exc:
372 raise HTTPException(
373 status.HTTP_409_CONFLICT,
374 detail=f"Error updating deployment: {exc}",
375 )
376 except CircularSchemaRefError:
377 raise HTTPException(
378 status_code=status.HTTP_422_UNPROCESSABLE_ENTITY,
379 detail="Invalid schema: Unable to validate schema with circular references.",
380 )
382 if deployment.global_concurrency_limit_id:
383 concurrency_limit = (
384 await models.concurrency_limits_v2.read_concurrency_limit(
385 session=session,
386 concurrency_limit_id=deployment.global_concurrency_limit_id,
387 )
388 )
390 if not concurrency_limit:
391 raise HTTPException(
392 status_code=status.HTTP_404_NOT_FOUND,
393 detail="Concurrency limit not found",
394 )
396 result = await models.deployments.update_deployment(
397 session=session,
398 deployment_id=deployment_id,
399 deployment=deployment,
400 )
402 # Phase 1: For schedules with `replaces`, look up the row ID by
403 # old slug, then NULL out those slugs so the unique index on
404 # (deployment_id, slug) won't block the renames.
405 renamed_schedule_ids: dict[str, UUID] = {}
406 renamed_old_slugs: list[str] = []
407 for schedule in schedules_to_patch:
408 if schedule.replaces:
409 row = (
410 await session.execute(
411 sa.select(db.DeploymentSchedule.id).where(
412 sa.and_(
413 db.DeploymentSchedule.deployment_id == deployment_id,
414 db.DeploymentSchedule.slug == schedule.replaces,
415 )
416 )
417 )
418 ).scalar_one_or_none()
419 if row is None:
420 raise HTTPException(
421 status_code=status.HTTP_404_NOT_FOUND,
422 detail=f"Schedule with slug '{schedule.replaces}' no longer exists.",
423 )
424 renamed_schedule_ids[schedule.replaces] = row
425 renamed_old_slugs.append(schedule.replaces)
427 if renamed_old_slugs:
428 await session.execute(
429 sa.update(db.DeploymentSchedule)
430 .where(
431 sa.and_(
432 db.DeploymentSchedule.deployment_id == deployment_id,
433 db.DeploymentSchedule.slug.in_(renamed_old_slugs),
434 )
435 )
436 .values(slug=None)
437 )
438 await session.flush()
440 # Phase 2: Apply the actual updates. Renamed schedules use
441 # their row ID (slug is now NULL); others use slug as before.
442 for schedule in schedules_to_patch:
443 if schedule.replaces:
444 await models.deployments.update_deployment_schedule(
445 session=session,
446 deployment_id=deployment_id,
447 schedule=schedule,
448 deployment_schedule_id=renamed_schedule_ids[schedule.replaces],
449 )
450 else:
451 await models.deployments.update_deployment_schedule(
452 session=session,
453 deployment_id=deployment_id,
454 schedule=schedule,
455 deployment_schedule_slug=schedule.slug,
456 )
457 if schedules_to_create:
458 await models.deployments.create_deployment_schedules(
459 session=session,
460 deployment_id=deployment_id,
461 schedules=[
462 schemas.actions.DeploymentScheduleCreate(
463 schedule=schedule.schedule, # type: ignore We will raise above if schedule is not provided
464 active=schedule.active if schedule.active is not None else True,
465 slug=schedule.slug,
466 parameters=schedule.parameters,
467 )
468 for schedule in schedules_to_create
469 ],
470 )
471 if not result:
472 raise HTTPException(status.HTTP_404_NOT_FOUND, detail="Deployment not found.")
475@router.get("/name/{flow_name}/{deployment_name}")
476async def read_deployment_by_name(
477 flow_name: str = Path(..., description="The name of the flow"),
478 deployment_name: str = Path(..., description="The name of the deployment"),
479 db: PrefectDBInterface = Depends(provide_database_interface),
480) -> schemas.responses.DeploymentResponse:
481 """
482 Get a deployment using the name of the flow and the deployment.
483 """
484 async with db.session_context() as session:
485 deployment = await models.deployments.read_deployment_by_name(
486 session=session, name=deployment_name, flow_name=flow_name
487 )
488 if not deployment:
489 raise HTTPException(
490 status.HTTP_404_NOT_FOUND, detail="Deployment not found"
491 )
492 return schemas.responses.DeploymentResponse.model_validate(
493 deployment, from_attributes=True
494 )
497@router.get("/{id:uuid}")
498async def read_deployment(
499 deployment_id: UUID = Path(..., description="The deployment id", alias="id"),
500 db: PrefectDBInterface = Depends(provide_database_interface),
501) -> schemas.responses.DeploymentResponse:
502 """
503 Get a deployment by id.
504 """
505 async with db.session_context() as session:
506 deployment = await models.deployments.read_deployment(
507 session=session, deployment_id=deployment_id
508 )
509 if not deployment:
510 raise HTTPException(
511 status_code=status.HTTP_404_NOT_FOUND, detail="Deployment not found"
512 )
513 return schemas.responses.DeploymentResponse.model_validate(
514 deployment, from_attributes=True
515 )
518@router.post("/filter")
519async def read_deployments(
520 limit: int = dependencies.LimitBody(),
521 offset: int = Body(0, ge=0),
522 flows: Optional[schemas.filters.FlowFilter] = None,
523 flow_runs: Optional[schemas.filters.FlowRunFilter] = None,
524 task_runs: Optional[schemas.filters.TaskRunFilter] = None,
525 deployments: Optional[schemas.filters.DeploymentFilter] = None,
526 work_pools: Optional[schemas.filters.WorkPoolFilter] = None,
527 work_pool_queues: Optional[schemas.filters.WorkQueueFilter] = None,
528 sort: schemas.sorting.DeploymentSort = Body(
529 schemas.sorting.DeploymentSort.NAME_ASC
530 ),
531 db: PrefectDBInterface = Depends(provide_database_interface),
532) -> List[schemas.responses.DeploymentResponse]:
533 """
534 Query for deployments.
535 """
536 async with db.session_context() as session:
537 response = await models.deployments.read_deployments(
538 session=session,
539 offset=offset,
540 sort=sort,
541 limit=limit,
542 flow_filter=flows,
543 flow_run_filter=flow_runs,
544 task_run_filter=task_runs,
545 deployment_filter=deployments,
546 work_pool_filter=work_pools,
547 work_queue_filter=work_pool_queues,
548 )
549 return [
550 schemas.responses.DeploymentResponse.model_validate(
551 deployment, from_attributes=True
552 )
553 for deployment in response
554 ]
557@router.post("/paginate")
558async def paginate_deployments(
559 limit: int = dependencies.LimitBody(),
560 page: int = Body(1, ge=1),
561 flows: Optional[schemas.filters.FlowFilter] = None,
562 flow_runs: Optional[schemas.filters.FlowRunFilter] = None,
563 task_runs: Optional[schemas.filters.TaskRunFilter] = None,
564 deployments: Optional[schemas.filters.DeploymentFilter] = None,
565 work_pools: Optional[schemas.filters.WorkPoolFilter] = None,
566 work_pool_queues: Optional[schemas.filters.WorkQueueFilter] = None,
567 sort: schemas.sorting.DeploymentSort = Body(
568 schemas.sorting.DeploymentSort.NAME_ASC
569 ),
570 db: PrefectDBInterface = Depends(provide_database_interface),
571) -> DeploymentPaginationResponse:
572 """
573 Pagination query for flow runs.
574 """
575 offset = (page - 1) * limit
577 async with db.session_context() as session:
578 response = await models.deployments.read_deployments(
579 session=session,
580 offset=offset,
581 sort=sort,
582 limit=limit,
583 flow_filter=flows,
584 flow_run_filter=flow_runs,
585 task_run_filter=task_runs,
586 deployment_filter=deployments,
587 work_pool_filter=work_pools,
588 work_queue_filter=work_pool_queues,
589 )
591 count = await models.deployments.count_deployments(
592 session=session,
593 flow_filter=flows,
594 flow_run_filter=flow_runs,
595 task_run_filter=task_runs,
596 deployment_filter=deployments,
597 work_pool_filter=work_pools,
598 work_queue_filter=work_pool_queues,
599 )
601 results = [
602 schemas.responses.DeploymentResponse.model_validate(
603 deployment, from_attributes=True
604 )
605 for deployment in response
606 ]
608 return DeploymentPaginationResponse(
609 results=results,
610 count=count,
611 limit=limit,
612 pages=(count + limit - 1) // limit,
613 page=page,
614 )
617@router.post("/get_scheduled_flow_runs")
618async def get_scheduled_flow_runs_for_deployments(
619 docket: dependencies.Docket,
620 deployment_ids: list[UUID] = Body(
621 default=..., description="The deployment IDs to get scheduled runs for"
622 ),
623 scheduled_before: DateTime = Body(
624 None, description="The maximum time to look for scheduled flow runs"
625 ),
626 limit: int = dependencies.LimitBody(),
627 db: PrefectDBInterface = Depends(provide_database_interface),
628) -> list[schemas.responses.FlowRunResponse]:
629 """
630 Get scheduled runs for a set of deployments. Used by a runner to poll for work.
631 """
632 async with db.session_context() as session:
633 orm_flow_runs = await models.flow_runs.read_flow_runs(
634 session=session,
635 limit=limit,
636 deployment_filter=schemas.filters.DeploymentFilter(
637 id=schemas.filters.DeploymentFilterId(any_=deployment_ids),
638 ),
639 flow_run_filter=schemas.filters.FlowRunFilter(
640 next_scheduled_start_time=schemas.filters.FlowRunFilterNextScheduledStartTime(
641 before_=scheduled_before
642 ),
643 state=schemas.filters.FlowRunFilterState(
644 type=schemas.filters.FlowRunFilterStateType(
645 any_=[schemas.states.StateType.SCHEDULED]
646 )
647 ),
648 ),
649 sort=schemas.sorting.FlowRunSort.NEXT_SCHEDULED_START_TIME_ASC,
650 )
652 flow_run_responses = [
653 schemas.responses.FlowRunResponse.model_validate(
654 orm_flow_run, from_attributes=True
655 )
656 for orm_flow_run in orm_flow_runs
657 ]
659 sorted_deployment_ids = ",".join(str(d) for d in sorted(deployment_ids))
660 await docket.add(
661 mark_deployments_ready,
662 key=f"mark_deployments_ready:deployments:{sorted_deployment_ids}",
663 )(
664 deployment_ids=deployment_ids,
665 )
667 return flow_run_responses
670@router.post("/count")
671async def count_deployments(
672 flows: Optional[schemas.filters.FlowFilter] = None,
673 flow_runs: Optional[schemas.filters.FlowRunFilter] = None,
674 task_runs: Optional[schemas.filters.TaskRunFilter] = None,
675 deployments: Optional[schemas.filters.DeploymentFilter] = None,
676 work_pools: Optional[schemas.filters.WorkPoolFilter] = None,
677 work_pool_queues: Optional[schemas.filters.WorkQueueFilter] = None,
678 db: PrefectDBInterface = Depends(provide_database_interface),
679) -> int:
680 """
681 Count deployments.
682 """
683 async with db.session_context() as session:
684 return await models.deployments.count_deployments(
685 session=session,
686 flow_filter=flows,
687 flow_run_filter=flow_runs,
688 task_run_filter=task_runs,
689 deployment_filter=deployments,
690 work_pool_filter=work_pools,
691 work_queue_filter=work_pool_queues,
692 )
695@router.delete("/{id:uuid}", status_code=status.HTTP_204_NO_CONTENT)
696async def delete_deployment(
697 deployment_id: UUID = Path(..., description="The deployment id", alias="id"),
698 db: PrefectDBInterface = Depends(provide_database_interface),
699) -> None:
700 """
701 Delete a deployment by id.
702 """
703 async with db.session_context(begin_transaction=True) as session:
704 result = await models.deployments.delete_deployment(
705 session=session, deployment_id=deployment_id
706 )
707 if not result: 707 ↛ exitline 707 didn't return from function 'delete_deployment' because the condition on line 707 was always true
708 raise HTTPException(
709 status_code=status.HTTP_404_NOT_FOUND, detail="Deployment not found"
710 )
713BULK_OPERATION_LIMIT = 50
716@router.post("/bulk_delete")
717async def bulk_delete_deployments(
718 deployments: Optional[schemas.filters.DeploymentFilter] = Body(
719 None, description="Filter criteria for deployments to delete"
720 ),
721 limit: int = Body(
722 BULK_OPERATION_LIMIT,
723 ge=1,
724 le=BULK_OPERATION_LIMIT,
725 description=f"Maximum number of deployments to delete. Defaults to {BULK_OPERATION_LIMIT}.",
726 ),
727 db: PrefectDBInterface = Depends(provide_database_interface),
728) -> DeploymentBulkDeleteResponse:
729 """
730 Bulk delete deployments matching the specified filter criteria.
732 Returns the IDs of deployments that were deleted.
733 """
734 async with db.session_context(begin_transaction=True) as session:
735 # Query matching deployments
736 db_deployments = await models.deployments.read_deployments(
737 session=session,
738 deployment_filter=deployments,
739 limit=limit,
740 )
742 if not db_deployments:
743 return DeploymentBulkDeleteResponse(deleted=[])
745 deployment_ids = [d.id for d in db_deployments]
747 # Delete deployments
748 deleted_ids = await models.deployments.delete_deployments(
749 session=session,
750 deployment_ids=deployment_ids,
751 )
753 return DeploymentBulkDeleteResponse(deleted=deleted_ids)
756@router.post("/{id:uuid}/schedule")
757async def schedule_deployment(
758 deployment_id: UUID = Path(..., description="The deployment id", alias="id"),
759 start_time: datetime.datetime = Body(
760 None, description="The earliest date to schedule"
761 ),
762 end_time: datetime.datetime = Body(None, description="The latest date to schedule"),
763 # Workaround for the fact that FastAPI does not let us configure ser_json_timedelta
764 # to represent timedeltas as floats in JSON.
765 min_time: float = Body(
766 None,
767 description=(
768 "Runs will be scheduled until at least this long after the `start_time`"
769 ),
770 json_schema_extra={"format": "time-delta"},
771 ),
772 min_runs: int = Body(None, description="The minimum number of runs to schedule"),
773 max_runs: int = Body(None, description="The maximum number of runs to schedule"),
774 db: PrefectDBInterface = Depends(provide_database_interface),
775) -> None:
776 """
777 Schedule runs for a deployment. For backfills, provide start/end times in the past.
779 This function will generate the minimum number of runs that satisfy the min
780 and max times, and the min and max counts. Specifically, the following order
781 will be respected.
783 - Runs will be generated starting on or after the `start_time`
784 - No more than `max_runs` runs will be generated
785 - No runs will be generated after `end_time` is reached
786 - At least `min_runs` runs will be generated
787 - Runs will be generated until at least `start_time + min_time` is reached
788 """
789 if isinstance(min_time, float):
790 min_time = datetime.timedelta(seconds=min_time)
792 async with db.session_context(begin_transaction=True) as session:
793 await models.deployments.schedule_runs(
794 session=session,
795 deployment_id=deployment_id,
796 start_time=start_time,
797 min_time=min_time,
798 end_time=end_time,
799 min_runs=min_runs,
800 max_runs=max_runs,
801 )
804@router.post("/{id:uuid}/resume_deployment")
805async def resume_deployment(
806 deployment_id: UUID = Path(..., description="The deployment id", alias="id"),
807 db: PrefectDBInterface = Depends(provide_database_interface),
808) -> None:
809 """
810 Set a deployment schedule to active. Runs will be scheduled immediately.
811 """
812 async with db.session_context(begin_transaction=True) as session:
813 deployment = await models.deployments.read_deployment(
814 session=session, deployment_id=deployment_id
815 )
816 if not deployment:
817 raise HTTPException(
818 status_code=status.HTTP_404_NOT_FOUND, detail="Deployment not found"
819 )
820 deployment.paused = False
823@router.post("/{id:uuid}/pause_deployment")
824async def pause_deployment(
825 deployment_id: UUID = Path(..., description="The deployment id", alias="id"),
826 db: PrefectDBInterface = Depends(provide_database_interface),
827) -> None:
828 """
829 Set a deployment schedule to inactive. Any auto-scheduled runs still in a Scheduled
830 state will be deleted.
831 """
832 async with db.session_context(begin_transaction=False) as session:
833 deployment = await models.deployments.read_deployment(
834 session=session, deployment_id=deployment_id
835 )
836 if not deployment:
837 raise HTTPException(
838 status_code=status.HTTP_404_NOT_FOUND, detail="Deployment not found"
839 )
840 deployment.paused = True
842 # commit here to make the inactive schedule "visible" to the scheduler service
843 await session.commit()
845 # delete any auto scheduled runs
846 await models.deployments._delete_scheduled_runs(
847 session=session,
848 deployment_id=deployment_id,
849 auto_scheduled_only=True,
850 )
852 await session.commit()
855@router.post("/{id:uuid}/create_flow_run")
856async def create_flow_run_from_deployment(
857 flow_run: schemas.actions.DeploymentFlowRunCreate,
858 deployment_id: UUID = Path(..., description="The deployment id", alias="id"),
859 created_by: Optional[schemas.core.CreatedBy] = Depends(dependencies.get_created_by),
860 db: PrefectDBInterface = Depends(provide_database_interface),
861 worker_lookups: WorkerLookups = Depends(WorkerLookups),
862 response: Response = None,
863) -> schemas.responses.FlowRunResponse:
864 """
865 Create a flow run from a deployment.
867 Any parameters not provided will be inferred from the deployment's parameters.
868 If tags are not provided, the deployment's tags will be used.
870 If no state is provided, the flow run will be created in a SCHEDULED state.
871 """
872 async with db.session_context(begin_transaction=True) as session:
873 # get relevant info from the deployment
874 deployment = await models.deployments.read_deployment(
875 session=session, deployment_id=deployment_id
876 )
878 if not deployment:
879 raise HTTPException(
880 status_code=status.HTTP_404_NOT_FOUND, detail="Deployment not found"
881 )
883 try:
884 dehydrated_params = deployment.parameters
885 dehydrated_params.update(flow_run.parameters or {})
886 ctx = await HydrationContext.build(
887 session=session,
888 raise_on_error=True,
889 render_jinja=True,
890 render_workspace_variables=True,
891 )
892 parameters = hydrate(dehydrated_params, ctx)
893 except HydrationError as exc:
894 raise HTTPException(
895 status.HTTP_400_BAD_REQUEST,
896 detail=f"Error hydrating flow run parameters: {exc}",
897 )
899 # default
900 enforce_parameter_schema = deployment.enforce_parameter_schema
902 # run override
903 if flow_run.enforce_parameter_schema is not None:
904 enforce_parameter_schema = flow_run.enforce_parameter_schema
906 if enforce_parameter_schema:
907 if not isinstance(deployment.parameter_openapi_schema, dict):
908 raise HTTPException(
909 status.HTTP_409_CONFLICT,
910 detail=(
911 "Error updating deployment: Cannot update parameters because"
912 " parameter schema enforcement is enabled and the deployment"
913 " does not have a valid parameter schema."
914 ),
915 )
916 try:
917 validate(
918 parameters, deployment.parameter_openapi_schema, raise_on_error=True
919 )
920 except ValidationError as exc:
921 raise HTTPException(
922 status.HTTP_409_CONFLICT,
923 detail=f"Error creating flow run: {exc}",
924 )
925 except CircularSchemaRefError:
926 raise HTTPException(
927 status_code=status.HTTP_422_UNPROCESSABLE_ENTITY,
928 detail="Invalid schema: Unable to validate schema with circular references.",
929 )
931 await validate_job_variables_for_deployment_flow_run(
932 session, deployment, flow_run
933 )
935 work_queue_name = deployment.work_queue_name
936 work_queue_id = deployment.work_queue_id
938 if flow_run.work_queue_name:
939 # can't mutate the ORM model or else it will commit the changes back
940 if deployment.work_queue is None or deployment.work_queue.work_pool is None:
941 raise HTTPException(
942 status_code=status.HTTP_400_BAD_REQUEST,
943 detail=f"Cannot create flow run in work queue {flow_run.work_queue_name} because deployment {deployment_id} is not associated with a work pool. Please remove work_pool_name and try again.",
944 )
946 work_queue_id = await worker_lookups._get_work_queue_id_from_name(
947 session=session,
948 work_pool_name=deployment.work_queue.work_pool.name,
949 work_queue_name=flow_run.work_queue_name,
950 create_queue_if_not_found=True,
951 )
952 work_queue_name = flow_run.work_queue_name
954 # hydrate the input model into a full flow run / state model
955 flow_run = schemas.core.FlowRun(
956 **flow_run.model_dump(
957 exclude={
958 "parameters",
959 "tags",
960 "infrastructure_document_id",
961 "work_queue_name",
962 "enforce_parameter_schema",
963 }
964 ),
965 flow_id=deployment.flow_id,
966 deployment_id=deployment.id,
967 deployment_version=deployment.version,
968 parameters=parameters,
969 tags=set(deployment.tags).union(flow_run.tags),
970 infrastructure_document_id=(
971 flow_run.infrastructure_document_id
972 or deployment.infrastructure_document_id
973 ),
974 work_queue_name=work_queue_name,
975 work_queue_id=work_queue_id,
976 created_by=created_by,
977 )
979 if not flow_run.state:
980 flow_run.state = schemas.states.Scheduled()
982 right_now = now("UTC")
983 model = await models.flow_runs.create_flow_run(
984 session=session, flow_run=flow_run
985 )
986 if model.created >= right_now:
987 response.status_code = status.HTTP_201_CREATED
988 return schemas.responses.FlowRunResponse.model_validate(
989 model, from_attributes=True
990 )
993BULK_CREATE_LIMIT = 100
996@router.post("/{id:uuid}/create_flow_run/bulk")
997async def bulk_create_flow_runs_from_deployment(
998 flow_runs: List[schemas.actions.DeploymentFlowRunCreate] = Body(
999 ..., description="List of flow run configurations to create"
1000 ),
1001 deployment_id: UUID = Path(..., description="The deployment id", alias="id"),
1002 created_by: Optional[schemas.core.CreatedBy] = Depends(dependencies.get_created_by),
1003 db: PrefectDBInterface = Depends(provide_database_interface),
1004 worker_lookups: WorkerLookups = Depends(WorkerLookups),
1005) -> FlowRunBulkCreateResponse:
1006 """
1007 Create multiple flow runs from a deployment.
1009 Any parameters not provided will be inferred from the deployment's parameters.
1010 If tags are not provided, the deployment's tags will be used.
1012 If no state is provided, the flow runs will be created in a SCHEDULED state.
1013 """
1014 if len(flow_runs) > BULK_CREATE_LIMIT: 1014 ↛ 1015line 1014 didn't jump to line 1015 because the condition on line 1014 was never true
1015 raise HTTPException(
1016 status_code=status.HTTP_400_BAD_REQUEST,
1017 detail=f"Cannot create more than {BULK_CREATE_LIMIT} flow runs at once.",
1018 )
1020 results: List[FlowRunCreateResult] = []
1022 async with db.session_context(begin_transaction=True) as session:
1023 # Get the deployment once - do this before the empty check so we
1024 # return 404 for non-existent deployments even with an empty list
1025 deployment = await models.deployments.read_deployment(
1026 session=session, deployment_id=deployment_id
1027 )
1029 if not deployment:
1030 raise HTTPException(
1031 status_code=status.HTTP_404_NOT_FOUND, detail="Deployment not found"
1032 )
1034 # Return early for empty list, but only after verifying deployment exists
1035 if not flow_runs:
1036 return FlowRunBulkCreateResponse(results=[])
1038 # Pre-create unique work queues to avoid race conditions
1039 # Collect unique work queue names that differ from the deployment's default
1040 unique_work_queue_names = {
1041 fr.work_queue_name
1042 for fr in flow_runs
1043 if fr.work_queue_name and fr.work_queue_name != deployment.work_queue_name
1044 }
1046 # Pre-create work queues if needed
1047 if (
1048 unique_work_queue_names
1049 and deployment.work_queue
1050 and deployment.work_queue.work_pool
1051 ):
1052 work_pool_name = deployment.work_queue.work_pool.name
1053 for work_queue_name in unique_work_queue_names:
1054 await worker_lookups._get_work_queue_id_from_name(
1055 session=session,
1056 work_pool_name=work_pool_name,
1057 work_queue_name=work_queue_name,
1058 create_queue_if_not_found=True,
1059 )
1061 # Build hydration context once
1062 try:
1063 ctx = await HydrationContext.build(
1064 session=session,
1065 raise_on_error=True,
1066 render_jinja=True,
1067 render_workspace_variables=True,
1068 )
1069 except HydrationError as exc:
1070 raise HTTPException(
1071 status.HTTP_400_BAD_REQUEST,
1072 detail=f"Error building hydration context: {exc}",
1073 )
1075 # Process flow runs sequentially within the transaction
1076 # (SQLAlchemy sessions are not safe for concurrent operations)
1077 for flow_run_request in flow_runs:
1078 try:
1079 # Hydrate parameters
1080 dehydrated_params = deployment.parameters.copy()
1081 dehydrated_params.update(flow_run_request.parameters or {})
1082 parameters = hydrate(dehydrated_params, ctx)
1084 # Default and override for enforce_parameter_schema
1085 enforce_parameter_schema = deployment.enforce_parameter_schema
1086 if flow_run_request.enforce_parameter_schema is not None:
1087 enforce_parameter_schema = flow_run_request.enforce_parameter_schema
1089 # Validate parameters if schema enforcement is enabled
1090 if enforce_parameter_schema:
1091 if not isinstance(deployment.parameter_openapi_schema, dict):
1092 results.append(
1093 FlowRunCreateResult(
1094 status="FAILED",
1095 error="Parameter schema enforcement is enabled but deployment has no valid schema.",
1096 )
1097 )
1098 continue
1099 try:
1100 validate(
1101 parameters,
1102 deployment.parameter_openapi_schema,
1103 raise_on_error=True,
1104 )
1105 except ValidationError as exc:
1106 results.append(
1107 FlowRunCreateResult(
1108 status="FAILED",
1109 error=f"Parameter validation failed: {exc}",
1110 )
1111 )
1112 continue
1113 except CircularSchemaRefError:
1114 results.append(
1115 FlowRunCreateResult(
1116 status="FAILED",
1117 error="Invalid schema: circular references detected.",
1118 )
1119 )
1120 continue
1122 # Validate job variables
1123 try:
1124 await validate_job_variables_for_deployment_flow_run(
1125 session, deployment, flow_run_request
1126 )
1127 except HTTPException as exc:
1128 results.append(
1129 FlowRunCreateResult(
1130 status="FAILED",
1131 error=str(exc.detail),
1132 )
1133 )
1134 continue
1136 # Determine work queue
1137 work_queue_name = deployment.work_queue_name
1138 work_queue_id = deployment.work_queue_id
1140 if flow_run_request.work_queue_name:
1141 if (
1142 deployment.work_queue is None
1143 or deployment.work_queue.work_pool is None
1144 ):
1145 results.append(
1146 FlowRunCreateResult(
1147 status="FAILED",
1148 error=f"Cannot create flow run in work queue {flow_run_request.work_queue_name} because deployment is not associated with a work pool.",
1149 )
1150 )
1151 continue
1153 work_queue_id = await worker_lookups._get_work_queue_id_from_name(
1154 session=session,
1155 work_pool_name=deployment.work_queue.work_pool.name,
1156 work_queue_name=flow_run_request.work_queue_name,
1157 create_queue_if_not_found=True,
1158 )
1159 work_queue_name = flow_run_request.work_queue_name
1161 # Create the flow run model
1162 flow_run_model = schemas.core.FlowRun(
1163 **flow_run_request.model_dump(
1164 exclude={
1165 "parameters",
1166 "tags",
1167 "infrastructure_document_id",
1168 "work_queue_name",
1169 "enforce_parameter_schema",
1170 }
1171 ),
1172 flow_id=deployment.flow_id,
1173 deployment_id=deployment.id,
1174 deployment_version=deployment.version,
1175 parameters=parameters,
1176 tags=set(deployment.tags).union(flow_run_request.tags),
1177 infrastructure_document_id=(
1178 flow_run_request.infrastructure_document_id
1179 or deployment.infrastructure_document_id
1180 ),
1181 work_queue_name=work_queue_name,
1182 work_queue_id=work_queue_id,
1183 created_by=created_by,
1184 )
1186 if not flow_run_model.state:
1187 flow_run_model.state = schemas.states.Scheduled()
1189 model = await models.flow_runs.create_flow_run(
1190 session=session, flow_run=flow_run_model
1191 )
1193 results.append(
1194 FlowRunCreateResult(
1195 flow_run_id=model.id,
1196 status="CREATED",
1197 )
1198 )
1200 except Exception as exc:
1201 results.append(
1202 FlowRunCreateResult(
1203 status="FAILED",
1204 error=str(exc),
1205 )
1206 )
1208 return FlowRunBulkCreateResponse(results=results)
1211# DEPRECATED
1212@router.get("/{id:uuid}/work_queue_check", deprecated=True)
1213async def work_queue_check_for_deployment(
1214 deployment_id: UUID = Path(..., description="The deployment id", alias="id"),
1215 db: PrefectDBInterface = Depends(provide_database_interface),
1216) -> List[schemas.core.WorkQueue]:
1217 """
1218 Get list of work-queues that are able to pick up the specified deployment.
1220 This endpoint is intended to be used by the UI to provide users warnings
1221 about deployments that are unable to be executed because there are no work
1222 queues that will pick up their runs, based on existing filter criteria. It
1223 may be deprecated in the future because there is not a strict relationship
1224 between work queues and deployments.
1225 """
1226 try:
1227 async with db.session_context() as session:
1228 work_queues = await models.deployments.check_work_queues_for_deployment(
1229 session=session, deployment_id=deployment_id
1230 )
1231 except ObjectNotFoundError:
1232 raise HTTPException(
1233 status_code=status.HTTP_404_NOT_FOUND, detail="Deployment not found"
1234 )
1235 return work_queues
1238@router.get("/{id:uuid}/schedules")
1239async def read_deployment_schedules(
1240 deployment_id: UUID = Path(..., description="The deployment id", alias="id"),
1241 db: PrefectDBInterface = Depends(provide_database_interface),
1242) -> List[schemas.core.DeploymentSchedule]:
1243 async with db.session_context() as session:
1244 deployment = await models.deployments.read_deployment(
1245 session=session, deployment_id=deployment_id
1246 )
1248 if not deployment:
1249 raise HTTPException(
1250 status.HTTP_404_NOT_FOUND, detail="Deployment not found."
1251 )
1253 return await models.deployments.read_deployment_schedules(
1254 session=session,
1255 deployment_id=deployment.id,
1256 )
1259@router.post("/{id:uuid}/schedules", status_code=status.HTTP_201_CREATED)
1260async def create_deployment_schedules(
1261 deployment_id: UUID = Path(..., description="The deployment id", alias="id"),
1262 schedules: List[schemas.actions.DeploymentScheduleCreate] = Body(
1263 default=..., description="The schedules to create"
1264 ),
1265 db: PrefectDBInterface = Depends(provide_database_interface),
1266) -> List[schemas.core.DeploymentSchedule]:
1267 async with db.session_context(begin_transaction=True) as session:
1268 deployment = await models.deployments.read_deployment(
1269 session=session, deployment_id=deployment_id
1270 )
1272 if not deployment:
1273 raise HTTPException(
1274 status.HTTP_404_NOT_FOUND, detail="Deployment not found."
1275 )
1277 try:
1278 created = await models.deployments.create_deployment_schedules(
1279 session=session,
1280 deployment_id=deployment.id,
1281 schedules=schedules,
1282 )
1283 except sa.exc.IntegrityError as e:
1284 if "duplicate key value violates unique constraint" in str(e):
1285 raise HTTPException(
1286 status.HTTP_409_CONFLICT,
1287 detail="Schedule slugs must be unique within a deployment.",
1288 )
1289 raise
1290 return created
1293@router.patch(
1294 "/{id:uuid}/schedules/{schedule_id:uuid}", status_code=status.HTTP_204_NO_CONTENT
1295)
1296async def update_deployment_schedule(
1297 deployment_id: UUID = Path(..., description="The deployment id", alias="id"),
1298 schedule_id: UUID = Path(..., description="The schedule id", alias="schedule_id"),
1299 schedule: schemas.actions.DeploymentScheduleUpdate = Body(
1300 default=..., description="The updated schedule"
1301 ),
1302 db: PrefectDBInterface = Depends(provide_database_interface),
1303) -> None:
1304 async with db.session_context(begin_transaction=True) as session:
1305 deployment = await models.deployments.read_deployment(
1306 session=session, deployment_id=deployment_id
1307 )
1309 if not deployment:
1310 raise HTTPException(
1311 status.HTTP_404_NOT_FOUND, detail="Deployment not found."
1312 )
1314 updated = await models.deployments.update_deployment_schedule(
1315 session=session,
1316 deployment_id=deployment_id,
1317 deployment_schedule_id=schedule_id,
1318 schedule=schedule,
1319 )
1321 if not updated:
1322 raise HTTPException(status.HTTP_404_NOT_FOUND, detail="Schedule not found.")
1324 await models.deployments._delete_scheduled_runs(
1325 session=session,
1326 deployment_id=deployment_id,
1327 auto_scheduled_only=True,
1328 future_only=True,
1329 )
1332@router.delete(
1333 "/{id:uuid}/schedules/{schedule_id:uuid}", status_code=status.HTTP_204_NO_CONTENT
1334)
1335async def delete_deployment_schedule(
1336 deployment_id: UUID = Path(..., description="The deployment id", alias="id"),
1337 schedule_id: UUID = Path(..., description="The schedule id", alias="schedule_id"),
1338 db: PrefectDBInterface = Depends(provide_database_interface),
1339) -> None:
1340 async with db.session_context(begin_transaction=True) as session:
1341 deployment = await models.deployments.read_deployment(
1342 session=session, deployment_id=deployment_id
1343 )
1345 if not deployment:
1346 raise HTTPException(
1347 status.HTTP_404_NOT_FOUND, detail="Deployment not found."
1348 )
1350 deleted = await models.deployments.delete_deployment_schedule(
1351 session=session,
1352 deployment_id=deployment_id,
1353 deployment_schedule_id=schedule_id,
1354 )
1356 if not deleted:
1357 raise HTTPException(status.HTTP_404_NOT_FOUND, detail="Schedule not found.")
1359 await models.deployments._delete_scheduled_runs(
1360 session=session,
1361 deployment_id=deployment_id,
1362 auto_scheduled_only=True,
1363 )