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

1""" 

2Routes for interacting with Deployment objects. 

3""" 

4 

5import datetime 

6import logging 

7from typing import List, Optional 

8from uuid import UUID 

9 

10import jsonschema.exceptions 

11import sqlalchemy as sa 

12from fastapi import Body, Depends, HTTPException, Path, Response 

13 

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) 

47 

48logger: logging.Logger = get_logger(__name__) 

49 

50router: PrefectRouter = PrefectRouter(prefix="/deployments", tags=["Deployments"]) 

51 

52 

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 ) 

63 

64 

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. 

77 

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. 

80 

81 For more information, see https://docs.prefect.io/v3/concepts/deployments. 

82 """ 

83 

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 

87 

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 ) 

102 

103 await validate_job_variables_for_deployment( 

104 session, 

105 work_pool, 

106 deployment, 

107 ) 

108 

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 ) 

114 

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 ) 

126 

127 if not concurrency_limit: 

128 raise HTTPException( 

129 status_code=status.HTTP_404_NOT_FOUND, 

130 detail="Concurrency limit not found", 

131 ) 

132 

133 deployment_dict["concurrency_limit_id"] = requested_concurrency_limit 

134 

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 

162 

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 ) 

183 

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 ) 

199 

200 right_now = now("UTC") 

201 model = await models.deployments.create_deployment( 

202 session=session, deployment=deployment 

203 ) 

204 

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 

207 

208 return schemas.responses.DeploymentResponse.model_validate( 

209 model, from_attributes=True 

210 ) 

211 

212 

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 ) 

227 

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 ] 

244 

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 "" 

255 

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 ) 

272 

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 

276 

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 ) 

306 

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 ) 

323 

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 

342 

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 

354 

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 ) 

381 

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 ) 

389 

390 if not concurrency_limit: 

391 raise HTTPException( 

392 status_code=status.HTTP_404_NOT_FOUND, 

393 detail="Concurrency limit not found", 

394 ) 

395 

396 result = await models.deployments.update_deployment( 

397 session=session, 

398 deployment_id=deployment_id, 

399 deployment=deployment, 

400 ) 

401 

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) 

426 

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

439 

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.") 

473 

474 

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 ) 

495 

496 

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 ) 

516 

517 

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 ] 

555 

556 

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 

576 

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 ) 

590 

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 ) 

600 

601 results = [ 

602 schemas.responses.DeploymentResponse.model_validate( 

603 deployment, from_attributes=True 

604 ) 

605 for deployment in response 

606 ] 

607 

608 return DeploymentPaginationResponse( 

609 results=results, 

610 count=count, 

611 limit=limit, 

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

613 page=page, 

614 ) 

615 

616 

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 ) 

651 

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 ] 

658 

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 ) 

666 

667 return flow_run_responses 

668 

669 

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 ) 

693 

694 

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 ) 

711 

712 

713BULK_OPERATION_LIMIT = 50 

714 

715 

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. 

731 

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 ) 

741 

742 if not db_deployments: 

743 return DeploymentBulkDeleteResponse(deleted=[]) 

744 

745 deployment_ids = [d.id for d in db_deployments] 

746 

747 # Delete deployments 

748 deleted_ids = await models.deployments.delete_deployments( 

749 session=session, 

750 deployment_ids=deployment_ids, 

751 ) 

752 

753 return DeploymentBulkDeleteResponse(deleted=deleted_ids) 

754 

755 

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. 

778 

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. 

782 

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) 

791 

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 ) 

802 

803 

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 

821 

822 

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 

841 

842 # commit here to make the inactive schedule "visible" to the scheduler service 

843 await session.commit() 

844 

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 ) 

851 

852 await session.commit() 

853 

854 

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. 

866 

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. 

869 

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 ) 

877 

878 if not deployment: 

879 raise HTTPException( 

880 status_code=status.HTTP_404_NOT_FOUND, detail="Deployment not found" 

881 ) 

882 

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 ) 

898 

899 # default 

900 enforce_parameter_schema = deployment.enforce_parameter_schema 

901 

902 # run override 

903 if flow_run.enforce_parameter_schema is not None: 

904 enforce_parameter_schema = flow_run.enforce_parameter_schema 

905 

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 ) 

930 

931 await validate_job_variables_for_deployment_flow_run( 

932 session, deployment, flow_run 

933 ) 

934 

935 work_queue_name = deployment.work_queue_name 

936 work_queue_id = deployment.work_queue_id 

937 

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 ) 

945 

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 

953 

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 ) 

978 

979 if not flow_run.state: 

980 flow_run.state = schemas.states.Scheduled() 

981 

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 ) 

991 

992 

993BULK_CREATE_LIMIT = 100 

994 

995 

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. 

1008 

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. 

1011 

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 ) 

1019 

1020 results: List[FlowRunCreateResult] = [] 

1021 

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 ) 

1028 

1029 if not deployment: 

1030 raise HTTPException( 

1031 status_code=status.HTTP_404_NOT_FOUND, detail="Deployment not found" 

1032 ) 

1033 

1034 # Return early for empty list, but only after verifying deployment exists 

1035 if not flow_runs: 

1036 return FlowRunBulkCreateResponse(results=[]) 

1037 

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 } 

1045 

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 ) 

1060 

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 ) 

1074 

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) 

1083 

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 

1088 

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 

1121 

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 

1135 

1136 # Determine work queue 

1137 work_queue_name = deployment.work_queue_name 

1138 work_queue_id = deployment.work_queue_id 

1139 

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 

1152 

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 

1160 

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 ) 

1185 

1186 if not flow_run_model.state: 

1187 flow_run_model.state = schemas.states.Scheduled() 

1188 

1189 model = await models.flow_runs.create_flow_run( 

1190 session=session, flow_run=flow_run_model 

1191 ) 

1192 

1193 results.append( 

1194 FlowRunCreateResult( 

1195 flow_run_id=model.id, 

1196 status="CREATED", 

1197 ) 

1198 ) 

1199 

1200 except Exception as exc: 

1201 results.append( 

1202 FlowRunCreateResult( 

1203 status="FAILED", 

1204 error=str(exc), 

1205 ) 

1206 ) 

1207 

1208 return FlowRunBulkCreateResponse(results=results) 

1209 

1210 

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. 

1219 

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 

1236 

1237 

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 ) 

1247 

1248 if not deployment: 

1249 raise HTTPException( 

1250 status.HTTP_404_NOT_FOUND, detail="Deployment not found." 

1251 ) 

1252 

1253 return await models.deployments.read_deployment_schedules( 

1254 session=session, 

1255 deployment_id=deployment.id, 

1256 ) 

1257 

1258 

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 ) 

1271 

1272 if not deployment: 

1273 raise HTTPException( 

1274 status.HTTP_404_NOT_FOUND, detail="Deployment not found." 

1275 ) 

1276 

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 

1291 

1292 

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 ) 

1308 

1309 if not deployment: 

1310 raise HTTPException( 

1311 status.HTTP_404_NOT_FOUND, detail="Deployment not found." 

1312 ) 

1313 

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 ) 

1320 

1321 if not updated: 

1322 raise HTTPException(status.HTTP_404_NOT_FOUND, detail="Schedule not found.") 

1323 

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 ) 

1330 

1331 

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 ) 

1344 

1345 if not deployment: 

1346 raise HTTPException( 

1347 status.HTTP_404_NOT_FOUND, detail="Deployment not found." 

1348 ) 

1349 

1350 deleted = await models.deployments.delete_deployment_schedule( 

1351 session=session, 

1352 deployment_id=deployment_id, 

1353 deployment_schedule_id=schedule_id, 

1354 ) 

1355 

1356 if not deleted: 

1357 raise HTTPException(status.HTTP_404_NOT_FOUND, detail="Schedule not found.") 

1358 

1359 await models.deployments._delete_scheduled_runs( 

1360 session=session, 

1361 deployment_id=deployment_id, 

1362 auto_scheduled_only=True, 

1363 )