Coverage for /usr/local/lib/python3.12/site-packages/prefect/server/models/deployments.py: 44%

386 statements  

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

1""" 

2Functions for interacting with deployment ORM objects. 

3Intended for internal use by the Prefect REST API. 

4""" 

5 

6from __future__ import annotations 

7 

8import datetime 

9import logging 

10from collections.abc import Iterable, Sequence 

11from typing import TYPE_CHECKING, Any, Optional, TypeVar, cast 

12from uuid import UUID 

13 

14import sqlalchemy as sa 

15from docket import Depends, Retry 

16from sqlalchemy import delete, or_, select 

17from sqlalchemy.dialects.postgresql import JSONB 

18from sqlalchemy.ext.asyncio import AsyncSession 

19from sqlalchemy.sql import Select 

20 

21from prefect._internal.uuid7 import uuid7 

22from prefect.logging import get_logger 

23from prefect.server import models, schemas 

24from prefect.server.database import ( 

25 PrefectDBInterface, 

26 db_injector, 

27 orm_models, 

28 provide_database_interface, 

29) 

30from prefect.server.events.clients import PrefectServerEventsClient 

31from prefect.server.exceptions import ObjectNotFoundError 

32from prefect.server.models.events import ( 

33 deployment_created_event, 

34 deployment_deleted_event, 

35 deployment_status_event, 

36 deployment_updated_event, 

37) 

38from prefect.server.schemas.statuses import DeploymentStatus 

39from prefect.settings import ( 

40 PREFECT_API_SERVICES_SCHEDULER_MAX_RUNS, 

41 PREFECT_API_SERVICES_SCHEDULER_MAX_SCHEDULED_TIME, 

42 PREFECT_API_SERVICES_SCHEDULER_MIN_RUNS, 

43 PREFECT_API_SERVICES_SCHEDULER_MIN_SCHEDULED_TIME, 

44) 

45from prefect.types._datetime import DateTime, now 

46 

47T = TypeVar("T", bound=tuple[Any, ...]) 

48 

49logger: logging.Logger = get_logger("prefect.server.models.deployments") 

50 

51DEPLOYMENT_EVENT_FIELDS = { 

52 "description", 

53 "tags", 

54 "parameters", 

55 "parameter_openapi_schema", 

56 "enforce_parameter_schema", 

57 "entrypoint", 

58 "path", 

59 "pull_steps", 

60 "work_queue_id", 

61 "infra_overrides", 

62 "paused", 

63 "labels", 

64 "version", 

65 "concurrency_limit_id", 

66 "concurrency_options", 

67} 

68 

69 

70@db_injector 

71async def _delete_scheduled_runs( 

72 db: PrefectDBInterface, 

73 session: AsyncSession, 

74 deployment_id: UUID, 

75 auto_scheduled_only: bool = False, 

76 future_only: bool = False, 

77) -> None: 

78 """ 

79 This utility function deletes all of a deployment's runs that are in a Scheduled state 

80 and haven't run yet. It should be run any time a deployment is created or 

81 modified in order to ensure that future runs comply with the deployment's latest values. 

82 

83 Args: 

84 deployment_id: the deployment for which we should delete runs. 

85 auto_scheduled_only: if True, only delete auto scheduled runs. Defaults to `False`. 

86 future_only: if True, only delete runs that are scheduled to run in the future. 

87 Defaults to `False`. 

88 """ 

89 delete_query = sa.delete(db.FlowRun).where( 

90 db.FlowRun.deployment_id == deployment_id, 

91 db.FlowRun.state_type == schemas.states.StateType.SCHEDULED.value, 

92 db.FlowRun.state_name != schemas.states.AwaitingConcurrencySlot().name, 

93 db.FlowRun.run_count == 0, 

94 ) 

95 

96 if auto_scheduled_only: 

97 delete_query = delete_query.where( 

98 db.FlowRun.auto_scheduled.is_(True), 

99 ) 

100 

101 if future_only: 

102 delete_query = delete_query.where( 

103 db.FlowRun.next_scheduled_start_time > now("UTC"), 

104 ) 

105 

106 await session.execute(delete_query) 

107 

108 

109@db_injector 

110async def create_deployment( 

111 db: PrefectDBInterface, 

112 session: AsyncSession, 

113 deployment: schemas.core.Deployment | schemas.actions.DeploymentCreate, 

114) -> Optional[orm_models.Deployment]: 

115 """Upserts a deployment. 

116 

117 Args: 

118 session: a database session 

119 deployment: a deployment model 

120 

121 Returns: 

122 orm_models.Deployment: the newly-created or updated deployment 

123 

124 """ 

125 

126 # Capture a timestamp before the upsert so we can compare against 

127 # result_deployment.created to reliably determine create-vs-update, 

128 # even under concurrent upserts for the same (flow_id, name). 

129 upsert_start = now("UTC") 

130 

131 # Snapshot existing deployment field values before upsert for change detection 

132 existing_result = await session.execute( 

133 sa.select(db.Deployment).where( 

134 sa.and_( 

135 db.Deployment.flow_id == deployment.flow_id, 

136 db.Deployment.name == deployment.name, 

137 ) 

138 ) 

139 ) 

140 existing_deployment = existing_result.scalar() 

141 existing_snapshot: Optional[dict[str, Any]] = None 

142 if existing_deployment is not None: 142 ↛ anywhereline 142 didn't jump anywhere: it always raised an exception.

143 existing_snapshot = { 

144 field: getattr(existing_deployment, field, None) 

145 for field in DEPLOYMENT_EVENT_FIELDS 

146 } 

147 # Expire to avoid stale ORM state after the Core-level upsert below 

148 session.expire(existing_deployment) 

149 

150 # set `updated` manually 

151 # known limitation of `on_conflict_do_update`, will not use `Column.onupdate` 

152 # https://docs.sqlalchemy.org/en/14/dialects/sqlite.html#the-set-clause 

153 deployment.updated = now("UTC") # type: ignore[assignment] 

154 

155 deployment.labels = await with_system_labels_for_deployment(session, deployment) 

156 

157 schedules = deployment.schedules 

158 insert_values = deployment.model_dump_for_orm( 

159 exclude_unset=True, exclude={"schedules", "version_info"} 

160 ) 

161 

162 requested_concurrency_limit = insert_values.pop("concurrency_limit", "unset") 

163 

164 # The job_variables field in client and server schemas is named 

165 # infra_overrides in the database. 

166 job_variables = insert_values.pop("job_variables", None) 

167 if job_variables: 167 ↛ anywhereline 167 didn't jump anywhere: it always raised an exception.

168 insert_values["infra_overrides"] = job_variables 

169 

170 conflict_update_fields = deployment.model_dump_for_orm( 

171 exclude_unset=True, 

172 exclude={ 

173 "id", 

174 "created", 

175 "created_by", 

176 "schedules", 

177 "job_variables", 

178 "concurrency_limit", 

179 "version_info", 

180 }, 

181 ) 

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

183 conflict_update_fields["infra_overrides"] = job_variables 

184 

185 insert_stmt = ( 

186 db.queries.insert(db.Deployment) 

187 .values(**insert_values) 

188 .on_conflict_do_update( 

189 index_elements=db.orm.deployment_unique_upsert_columns, 

190 set_={**conflict_update_fields}, 

191 ) 

192 ) 

193 

194 await session.execute(insert_stmt) 

195 

196 # Get the id of the deployment we just created or updated 

197 result = await session.execute( 

198 sa.select(db.Deployment.id).where( 

199 sa.and_( 

200 db.Deployment.flow_id == deployment.flow_id, 

201 db.Deployment.name == deployment.name, 

202 ) 

203 ) 

204 ) 

205 deployment_id = result.scalar_one_or_none() 

206 

207 if not deployment_id: 

208 return None 

209 

210 # Because this was possibly an upsert, we need to delete any existing 

211 # schedules and any runs from the old deployment. 

212 

213 await _delete_scheduled_runs( 

214 session=session, 

215 deployment_id=deployment_id, 

216 auto_scheduled_only=True, 

217 future_only=True, 

218 ) 

219 

220 await delete_schedules_for_deployment(session=session, deployment_id=deployment_id) 

221 

222 if schedules: 

223 await create_deployment_schedules( 

224 session=session, 

225 deployment_id=deployment_id, 

226 schedules=[ 

227 schemas.actions.DeploymentScheduleCreate( 

228 schedule=schedule.schedule, 

229 active=schedule.active, 

230 parameters=schedule.parameters, 

231 slug=schedule.slug, 

232 ) 

233 for schedule in schedules 

234 ], 

235 ) 

236 

237 if requested_concurrency_limit != "unset": 

238 await _create_or_update_deployment_concurrency_limit( 

239 db, session, deployment_id, deployment.concurrency_limit 

240 ) 

241 

242 query = ( 

243 sa.select(db.Deployment) 

244 .where( 

245 sa.and_( 

246 db.Deployment.flow_id == deployment.flow_id, 

247 db.Deployment.name == deployment.name, 

248 ) 

249 ) 

250 .execution_options(populate_existing=True) 

251 ) 

252 refreshed_result = await session.execute(query) 

253 result_deployment = refreshed_result.scalar() 

254 

255 if result_deployment is not None: 

256 if result_deployment.created >= upsert_start: 

257 # The row was genuinely inserted (not an ON CONFLICT update). 

258 await emit_deployment_created_event( 

259 session=session, deployment=result_deployment 

260 ) 

261 elif existing_snapshot is not None: 

262 changed_fields = _detect_deployment_changed_fields( 

263 existing_snapshot, result_deployment 

264 ) 

265 if changed_fields: 

266 await emit_deployment_updated_event( 

267 session=session, 

268 deployment=result_deployment, 

269 changed_fields=changed_fields, 

270 ) 

271 

272 return result_deployment 

273 

274 

275@db_injector 

276async def update_deployment( 

277 db: PrefectDBInterface, 

278 session: AsyncSession, 

279 deployment_id: UUID, 

280 deployment: schemas.actions.DeploymentUpdate, 

281) -> bool: 

282 """Updates a deployment. 

283 

284 Args: 

285 session: a database session 

286 deployment_id: the ID of the deployment to modify 

287 deployment: changes to a deployment model 

288 

289 Returns: 

290 bool: whether the deployment was updated 

291 

292 """ 

293 

294 from prefect.server.api.workers import WorkerLookups 

295 

296 # Snapshot current field values before update for change detection 

297 current_deployment = await read_deployment( 

298 session=session, deployment_id=deployment_id 

299 ) 

300 current_snapshot: Optional[dict[str, Any]] = None 

301 if current_deployment is not None: 

302 current_snapshot = { 

303 field: getattr(current_deployment, field, None) 

304 for field in DEPLOYMENT_EVENT_FIELDS 

305 } 

306 # Expire to avoid stale ORM state after the Core-level update below 

307 session.expire(current_deployment) 

308 

309 schedules = deployment.schedules 

310 

311 # exclude_unset=True allows us to only update values provided by 

312 # the user, ignoring any defaults on the model 

313 update_data = deployment.model_dump_for_orm( 

314 exclude_unset=True, 

315 exclude={"work_pool_name", "version_info"}, 

316 ) 

317 

318 requested_global_concurrency_limit_update = update_data.pop( 

319 "global_concurrency_limit_id", "unset" 

320 ) 

321 requested_concurrency_limit_update = update_data.pop("concurrency_limit", "unset") 

322 

323 if requested_global_concurrency_limit_update != "unset": 

324 update_data["concurrency_limit_id"] = requested_global_concurrency_limit_update 

325 

326 # The job_variables field in client and server schemas is named 

327 # infra_overrides in the database. 

328 job_variables = update_data.pop("job_variables", None) 

329 if job_variables: 

330 update_data["infra_overrides"] = job_variables 

331 

332 should_update_schedules = update_data.pop("schedules", None) is not None 

333 

334 if deployment.work_pool_name and deployment.work_queue_name: 

335 # If a specific pool name/queue name combination was provided, get the 

336 # ID for that work pool queue. 

337 update_data[ 

338 "work_queue_id" 

339 ] = await WorkerLookups()._get_work_queue_id_from_name( 

340 session=session, 

341 work_pool_name=deployment.work_pool_name, 

342 work_queue_name=deployment.work_queue_name, 

343 create_queue_if_not_found=True, 

344 ) 

345 elif deployment.work_pool_name: 

346 # If just a pool name was provided, get the ID for its default 

347 # work pool queue. 

348 update_data[ 

349 "work_queue_id" 

350 ] = await WorkerLookups()._get_default_work_queue_id_from_work_pool_name( 

351 session=session, 

352 work_pool_name=deployment.work_pool_name, 

353 ) 

354 elif ( 

355 deployment.work_pool_name is None 

356 and "work_pool_name" in deployment.model_fields_set 

357 ): 

358 # work_pool_name was explicitly set to None — clear the work queue so 

359 # runs are no longer routed to a work pool (e.g. switching to serve). 

360 update_data["work_queue_id"] = None 

361 elif deployment.work_queue_name: 

362 # If just a queue name was provided, ensure the queue exists and 

363 # get its ID. 

364 work_queue = await models.work_queues.ensure_work_queue_exists( 

365 session=session, name=update_data["work_queue_name"] 

366 ) 

367 update_data["work_queue_id"] = work_queue.id 

368 

369 update_stmt = ( 

370 sa.update(db.Deployment) 

371 .where(db.Deployment.id == deployment_id) 

372 .values(**update_data) 

373 ) 

374 result = await session.execute(update_stmt) 

375 

376 # delete any auto scheduled runs that would have reflected the old deployment config 

377 await _delete_scheduled_runs( 

378 session=session, 

379 deployment_id=deployment_id, 

380 auto_scheduled_only=True, 

381 future_only=True, 

382 ) 

383 

384 if should_update_schedules: 

385 # If schedules were provided, remove the existing schedules and 

386 # replace them with the new ones. 

387 await delete_schedules_for_deployment( 

388 session=session, deployment_id=deployment_id 

389 ) 

390 await create_deployment_schedules( 

391 session=session, 

392 deployment_id=deployment_id, 

393 schedules=[ 

394 schemas.actions.DeploymentScheduleCreate( 

395 schedule=schedule.schedule, 

396 active=schedule.active if schedule.active is not None else True, 

397 parameters=schedule.parameters, 

398 slug=schedule.slug, 

399 ) 

400 for schedule in schedules 

401 if schedule.schedule is not None 

402 ], 

403 ) 

404 

405 if requested_concurrency_limit_update != "unset": 

406 await _create_or_update_deployment_concurrency_limit( 

407 db, session, deployment_id, deployment.concurrency_limit 

408 ) 

409 

410 updated = result.rowcount > 0 

411 

412 if updated and current_snapshot is not None: 

413 updated_deployment = await read_deployment( 

414 session=session, deployment_id=deployment_id 

415 ) 

416 if updated_deployment is not None: 

417 changed_fields = _detect_deployment_changed_fields( 

418 current_snapshot, updated_deployment 

419 ) 

420 if changed_fields: 

421 await emit_deployment_updated_event( 

422 session=session, 

423 deployment=updated_deployment, 

424 changed_fields=changed_fields, 

425 ) 

426 

427 return updated 

428 

429 

430async def _create_or_update_deployment_concurrency_limit( 

431 db: PrefectDBInterface, 

432 session: AsyncSession, 

433 deployment_id: UUID, 

434 limit: Optional[int], 

435): 

436 deployment = await session.get(db.Deployment, deployment_id) 

437 assert deployment is not None 

438 

439 if ( 

440 deployment.global_concurrency_limit 

441 and deployment.global_concurrency_limit.limit == limit 

442 ) or (deployment.global_concurrency_limit is None and limit is None): 

443 return 

444 

445 deployment._concurrency_limit = limit 

446 if limit is None: 

447 await _delete_related_concurrency_limit( 

448 db, session=session, deployment_id=deployment_id 

449 ) 

450 await session.refresh(deployment) 

451 elif deployment.global_concurrency_limit: 

452 deployment.global_concurrency_limit.limit = limit 

453 else: 

454 limit_name = f"deployment:{deployment_id}" 

455 new_limit = db.ConcurrencyLimitV2(name=limit_name, limit=limit) 

456 deployment.global_concurrency_limit = new_limit 

457 

458 session.add(deployment) 

459 

460 

461@db_injector 

462async def read_deployment( 

463 db: PrefectDBInterface, session: AsyncSession, deployment_id: UUID 

464) -> Optional[orm_models.Deployment]: 

465 """Reads a deployment by id. 

466 

467 Args: 

468 session: A database session 

469 deployment_id: a deployment id 

470 

471 Returns: 

472 orm_models.Deployment: the deployment 

473 """ 

474 

475 return await session.get(db.Deployment, deployment_id) 

476 

477 

478@db_injector 

479async def read_deployment_by_name( 

480 db: PrefectDBInterface, session: AsyncSession, name: str, flow_name: str 

481) -> Optional[orm_models.Deployment]: 

482 """Reads a deployment by name. 

483 

484 Args: 

485 session: A database session 

486 name: a deployment name 

487 flow_name: the name of the flow the deployment belongs to 

488 

489 Returns: 

490 orm_models.Deployment: the deployment 

491 """ 

492 

493 result = await session.execute( 

494 select(db.Deployment) 

495 .join(db.Flow, db.Deployment.flow_id == db.Flow.id) 

496 .where( 

497 sa.and_( 

498 db.Flow.name == flow_name, 

499 db.Deployment.name == name, 

500 ) 

501 ) 

502 .limit(1) 

503 ) 

504 return result.scalar() 

505 

506 

507async def _apply_deployment_filters( 

508 db: PrefectDBInterface, 

509 query: Select[T], 

510 flow_filter: Optional[schemas.filters.FlowFilter] = None, 

511 flow_run_filter: Optional[schemas.filters.FlowRunFilter] = None, 

512 task_run_filter: Optional[schemas.filters.TaskRunFilter] = None, 

513 deployment_filter: Optional[schemas.filters.DeploymentFilter] = None, 

514 work_pool_filter: Optional[schemas.filters.WorkPoolFilter] = None, 

515 work_queue_filter: Optional[schemas.filters.WorkQueueFilter] = None, 

516) -> Select[T]: 

517 """ 

518 Applies filters to a deployment query as a combination of EXISTS subqueries. 

519 """ 

520 

521 if deployment_filter: 

522 query = query.where(deployment_filter.as_sql_filter()) 

523 

524 if flow_filter: 

525 flow_exists_clause = select(db.Deployment.id).where( 

526 db.Deployment.flow_id == db.Flow.id, 

527 flow_filter.as_sql_filter(), 

528 ) 

529 

530 query = query.where(flow_exists_clause.exists()) 

531 

532 if flow_run_filter or task_run_filter: 

533 flow_run_exists_clause = select(db.FlowRun).where( 

534 db.Deployment.id == db.FlowRun.deployment_id 

535 ) 

536 

537 if flow_run_filter: 

538 flow_run_exists_clause = flow_run_exists_clause.where( 

539 flow_run_filter.as_sql_filter() 

540 ) 

541 if task_run_filter: 

542 flow_run_exists_clause = flow_run_exists_clause.join( 

543 db.TaskRun, 

544 db.TaskRun.flow_run_id == db.FlowRun.id, 

545 ).where(task_run_filter.as_sql_filter()) 

546 

547 query = query.where(flow_run_exists_clause.exists()) 

548 

549 if work_pool_filter or work_queue_filter: 

550 work_pool_exists_clause = select(db.WorkQueue).where( 

551 db.Deployment.work_queue_id == db.WorkQueue.id 

552 ) 

553 

554 if work_queue_filter: 

555 work_pool_exists_clause = work_pool_exists_clause.where( 

556 work_queue_filter.as_sql_filter() 

557 ) 

558 

559 if work_pool_filter: 

560 work_pool_exists_clause = work_pool_exists_clause.join( 

561 db.WorkPool, 

562 db.WorkPool.id == db.WorkQueue.work_pool_id, 

563 ).where(work_pool_filter.as_sql_filter()) 

564 

565 query = query.where(work_pool_exists_clause.exists()) 

566 

567 return query 

568 

569 

570@db_injector 

571async def read_deployments( 

572 db: PrefectDBInterface, 

573 session: AsyncSession, 

574 offset: Optional[int] = None, 

575 limit: Optional[int] = None, 

576 flow_filter: Optional[schemas.filters.FlowFilter] = None, 

577 flow_run_filter: Optional[schemas.filters.FlowRunFilter] = None, 

578 task_run_filter: Optional[schemas.filters.TaskRunFilter] = None, 

579 deployment_filter: Optional[schemas.filters.DeploymentFilter] = None, 

580 work_pool_filter: Optional[schemas.filters.WorkPoolFilter] = None, 

581 work_queue_filter: Optional[schemas.filters.WorkQueueFilter] = None, 

582 sort: schemas.sorting.DeploymentSort = schemas.sorting.DeploymentSort.NAME_ASC, 

583) -> Sequence[orm_models.Deployment]: 

584 """ 

585 Read deployments. 

586 

587 Args: 

588 session: A database session 

589 offset: Query offset 

590 limit: Query limit 

591 flow_filter: only select deployments whose flows match these criteria 

592 flow_run_filter: only select deployments whose flow runs match these criteria 

593 task_run_filter: only select deployments whose task runs match these criteria 

594 deployment_filter: only select deployment that match these filters 

595 work_pool_filter: only select deployments whose work pools match these criteria 

596 work_queue_filter: only select deployments whose work pool queues match these criteria 

597 sort: the sort criteria for selected deployments. Defaults to `name` ASC. 

598 

599 Returns: 

600 list[orm_models.Deployment]: deployments 

601 """ 

602 

603 query = select(db.Deployment).order_by(*sort.as_sql_sort()) 

604 

605 query = await _apply_deployment_filters( 

606 db, 

607 query=query, 

608 flow_filter=flow_filter, 

609 flow_run_filter=flow_run_filter, 

610 task_run_filter=task_run_filter, 

611 deployment_filter=deployment_filter, 

612 work_pool_filter=work_pool_filter, 

613 work_queue_filter=work_queue_filter, 

614 ) 

615 

616 if offset is not None: 

617 query = query.offset(offset) 

618 if limit is not None: 618 ↛ 621line 618 didn't jump to line 621 because the condition on line 618 was always true

619 query = query.limit(limit) 

620 

621 result = await session.execute(query) 

622 return result.scalars().unique().all() 

623 

624 

625@db_injector 

626async def count_deployments( 

627 db: PrefectDBInterface, 

628 session: AsyncSession, 

629 flow_filter: Optional[schemas.filters.FlowFilter] = None, 

630 flow_run_filter: Optional[schemas.filters.FlowRunFilter] = None, 

631 task_run_filter: Optional[schemas.filters.TaskRunFilter] = None, 

632 deployment_filter: Optional[schemas.filters.DeploymentFilter] = None, 

633 work_pool_filter: Optional[schemas.filters.WorkPoolFilter] = None, 

634 work_queue_filter: Optional[schemas.filters.WorkQueueFilter] = None, 

635) -> int: 

636 """ 

637 Count deployments. 

638 

639 Args: 

640 session: A database session 

641 flow_filter: only count deployments whose flows match these criteria 

642 flow_run_filter: only count deployments whose flow runs match these criteria 

643 task_run_filter: only count deployments whose task runs match these criteria 

644 deployment_filter: only count deployment that match these filters 

645 work_pool_filter: only count deployments that match these work pool filters 

646 work_queue_filter: only count deployments that match these work pool queue filters 

647 

648 Returns: 

649 int: the number of deployments matching filters 

650 """ 

651 

652 query = select(sa.func.count(None)).select_from(db.Deployment) 

653 

654 query = await _apply_deployment_filters( 

655 db, 

656 query=query, 

657 flow_filter=flow_filter, 

658 flow_run_filter=flow_run_filter, 

659 task_run_filter=task_run_filter, 

660 deployment_filter=deployment_filter, 

661 work_pool_filter=work_pool_filter, 

662 work_queue_filter=work_queue_filter, 

663 ) 

664 

665 result = await session.execute(query) 

666 return result.scalar_one() 

667 

668 

669@db_injector 

670async def delete_deployment( 

671 db: PrefectDBInterface, session: AsyncSession, deployment_id: UUID 

672) -> bool: 

673 """ 

674 Delete a deployment by id. 

675 

676 Args: 

677 session: A database session 

678 deployment_id: a deployment id 

679 

680 Returns: 

681 bool: whether or not the deployment was deleted 

682 """ 

683 

684 # Build the delete event before deletion while the deployment is still in session 

685 deployment = await read_deployment(session=session, deployment_id=deployment_id) 

686 delete_event = None 

687 if deployment is not None: 

688 delete_event = await deployment_deleted_event( 

689 session=session, deployment=deployment, occurred=now("UTC") 

690 ) 

691 

692 # delete scheduled runs, both auto- and user- created. 

693 await _delete_scheduled_runs( 

694 session=session, deployment_id=deployment_id, auto_scheduled_only=False 

695 ) 

696 

697 await _delete_related_concurrency_limit( 

698 db, session=session, deployment_id=deployment_id 

699 ) 

700 

701 result = await session.execute( 

702 delete(db.Deployment).where(db.Deployment.id == deployment_id) 

703 ) 

704 deleted = result.rowcount > 0 

705 

706 if deleted and delete_event is not None: 

707 async with PrefectServerEventsClient() as events_client: 

708 await events_client.emit(delete_event) 

709 

710 return deleted 

711 

712 

713async def _delete_related_concurrency_limit( 

714 db: PrefectDBInterface, session: AsyncSession, deployment_id: UUID 

715): 

716 return await session.execute( 

717 delete(db.ConcurrencyLimitV2).where( 

718 db.ConcurrencyLimitV2.id 

719 == sa.select(db.Deployment.concurrency_limit_id) 

720 .where(db.Deployment.id == deployment_id) 

721 .scalar_subquery() 

722 ) 

723 ) 

724 

725 

726@db_injector 

727async def delete_deployments( 

728 db: PrefectDBInterface, 

729 session: AsyncSession, 

730 deployment_ids: list[UUID], 

731) -> list[UUID]: 

732 """ 

733 Delete multiple deployments by their IDs. 

734 

735 Args: 

736 session: A database session 

737 deployment_ids: a list of deployment ids to delete 

738 

739 Returns: 

740 List[UUID]: the IDs of the deployments that were deleted 

741 """ 

742 if not deployment_ids: 

743 return [] 

744 

745 # Build delete events before deletion while deployments are still in session 

746 result = await session.execute( 

747 select(db.Deployment).where(db.Deployment.id.in_(deployment_ids)) 

748 ) 

749 existing_deployments = list(result.scalars().unique().all()) 

750 

751 if not existing_deployments: 

752 return [] 

753 

754 existing_ids = [d.id for d in existing_deployments] 

755 

756 delete_events = [] 

757 for d in existing_deployments: 

758 delete_events.append( 

759 await deployment_deleted_event( 

760 session=session, deployment=d, occurred=now("UTC") 

761 ) 

762 ) 

763 

764 # Delete scheduled runs for all deployments 

765 for deployment_id in existing_ids: 

766 await _delete_scheduled_runs( 

767 session=session, deployment_id=deployment_id, auto_scheduled_only=False 

768 ) 

769 

770 # Delete related concurrency limits 

771 await session.execute( 

772 delete(db.ConcurrencyLimitV2).where( 

773 db.ConcurrencyLimitV2.id.in_( 

774 select(db.Deployment.concurrency_limit_id).where( 

775 db.Deployment.id.in_(existing_ids) 

776 ) 

777 ) 

778 ) 

779 ) 

780 

781 # Delete all deployments 

782 await session.execute( 

783 delete(db.Deployment).where(db.Deployment.id.in_(existing_ids)) 

784 ) 

785 

786 # Emit delete events 

787 async with PrefectServerEventsClient() as events_client: 

788 for event in delete_events: 

789 await events_client.emit(event) 

790 

791 return existing_ids 

792 

793 

794@db_injector 

795async def schedule_runs( 

796 db: PrefectDBInterface, 

797 session: AsyncSession, 

798 deployment_id: UUID, 

799 start_time: Optional[datetime.datetime] = None, 

800 end_time: Optional[datetime.datetime] = None, 

801 min_time: Optional[datetime.timedelta] = None, 

802 min_runs: Optional[int] = None, 

803 max_runs: Optional[int] = None, 

804 auto_scheduled: bool = True, 

805) -> Sequence[UUID]: 

806 """ 

807 Schedule flow runs for a deployment 

808 

809 Args: 

810 session: a database session 

811 deployment_id: the id of the deployment to schedule 

812 start_time: the time from which to start scheduling runs 

813 end_time: runs will be scheduled until at most this time 

814 min_time: runs will be scheduled until at least this far in the future 

815 min_runs: a minimum amount of runs to schedule 

816 max_runs: a maximum amount of runs to schedule 

817 

818 This function will generate the minimum number of runs that satisfy the min 

819 and max times, and the min and max counts. Specifically, the following order 

820 will be respected. 

821 

822 - Runs will be generated starting on or after the `start_time` 

823 - No more than `max_runs` runs will be generated 

824 - No runs will be generated after `end_time` is reached 

825 - At least `min_runs` runs will be generated 

826 - Runs will be generated until at least `start_time` + `min_time` is reached 

827 

828 Returns: 

829 a list of flow run ids scheduled for the deployment 

830 """ 

831 if min_runs is None: 

832 min_runs = PREFECT_API_SERVICES_SCHEDULER_MIN_RUNS.value() 

833 assert min_runs is not None 

834 if max_runs is None: 

835 max_runs = PREFECT_API_SERVICES_SCHEDULER_MAX_RUNS.value() 

836 assert max_runs is not None 

837 if start_time is None: 

838 start_time = now("UTC") 

839 if end_time is None: 

840 end_time = start_time + ( 

841 PREFECT_API_SERVICES_SCHEDULER_MAX_SCHEDULED_TIME.value() 

842 ) 

843 if min_time is None: 

844 min_time = PREFECT_API_SERVICES_SCHEDULER_MIN_SCHEDULED_TIME.value() 

845 assert min_time is not None 

846 

847 actual_start_time = start_time 

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

849 assert end_time is not None 

850 actual_end_time = end_time 

851 

852 runs = await _generate_scheduled_flow_runs( 

853 db, 

854 session=session, 

855 deployment_id=deployment_id, 

856 start_time=actual_start_time, 

857 end_time=actual_end_time, 

858 min_time=min_time, 

859 min_runs=min_runs, 

860 max_runs=max_runs, 

861 auto_scheduled=auto_scheduled, 

862 ) 

863 return await _insert_scheduled_flow_runs(session=session, runs=runs) 

864 

865 

866async def _generate_scheduled_flow_runs( 

867 db: PrefectDBInterface, 

868 session: AsyncSession, 

869 deployment_id: UUID, 

870 start_time: datetime.datetime, 

871 end_time: datetime.datetime, 

872 min_time: datetime.timedelta, 

873 min_runs: int, 

874 max_runs: int, 

875 auto_scheduled: bool = True, 

876) -> list[dict[str, Any]]: 

877 """ 

878 Given a `deployment_id` and schedule, generates a list of flow run objects and 

879 associated scheduled states that represent scheduled flow runs. This method 

880 does NOT insert generated runs into the database, in order to facilitate 

881 batch operations. Call `_insert_scheduled_flow_runs()` to insert these runs. 

882 

883 Runs include an idempotency key which prevents duplicate runs from being inserted 

884 if the output from this function is used more than once. 

885 

886 Args: 

887 session: a database session 

888 deployment_id: the id of the deployment to schedule 

889 start_time: the time from which to start scheduling runs 

890 end_time: runs will be scheduled until at most this time 

891 min_time: runs will be scheduled until at least this far in the future 

892 min_runs: a minimum amount of runs to schedule 

893 max_runs: a maximum amount of runs to schedule 

894 

895 This function will generate the minimum number of runs that satisfy the min 

896 and max times, and the min and max counts. Specifically, the following order 

897 will be respected. 

898 

899 - Runs will be generated starting on or after the `start_time` 

900 - No more than `max_runs` runs will be generated 

901 - No runs will be generated after `end_time` is reached 

902 - At least `min_runs` runs will be generated 

903 - Runs will be generated until at least `start_time + min_time` is reached 

904 

905 Returns: 

906 a list of dictionary representations of the `FlowRun` objects to schedule 

907 """ 

908 runs: list[dict[str, Any]] = [] 

909 

910 deployment = await session.get(db.Deployment, deployment_id) 

911 

912 if not deployment: 

913 return [] 

914 

915 active_deployment_schedules = await read_deployment_schedules( 

916 session=session, 

917 deployment_id=deployment.id, 

918 deployment_schedule_filter=schemas.filters.DeploymentScheduleFilter( 

919 active=schemas.filters.DeploymentScheduleFilterActive(eq_=True) 

920 ), 

921 ) 

922 

923 for deployment_schedule in active_deployment_schedules: 

924 dates: list[DateTime] = [] 

925 

926 # generate up to `n` dates satisfying the min of `max_runs` and `end_time` 

927 for dt in deployment_schedule.schedule._get_dates_generator( 

928 n=max_runs, start=start_time, end=end_time 

929 ): 

930 dates.append(dt) 

931 

932 # at any point, if we satisfy both of the minimums, we can stop 

933 if len(dates) >= min_runs and dt >= (start_time + min_time): 

934 break 

935 

936 tags = deployment.tags 

937 if auto_scheduled: 

938 tags = ["auto-scheduled"] + tags 

939 

940 parameters = { 

941 **deployment.parameters, 

942 **deployment_schedule.parameters, 

943 } 

944 

945 # Generate system labels for flow runs from this deployment 

946 labels = await with_system_labels_for_deployment_flow_run( 

947 session=session, 

948 deployment=deployment, 

949 ) 

950 

951 for date in dates: 

952 runs.append( 

953 { 

954 "id": uuid7(), 

955 "flow_id": deployment.flow_id, 

956 "deployment_id": deployment_id, 

957 "deployment_version": deployment.version, 

958 "work_queue_name": deployment.work_queue_name, 

959 "work_queue_id": deployment.work_queue_id, 

960 "parameters": parameters, 

961 "infrastructure_document_id": deployment.infrastructure_document_id, 

962 "idempotency_key": f"scheduled {deployment.id} {deployment_schedule.id} {date}", 

963 "tags": tags, 

964 "labels": labels, 

965 "auto_scheduled": auto_scheduled, 

966 "state": schemas.states.Scheduled( 

967 scheduled_time=date, 

968 message="Flow run scheduled", 

969 ).model_dump(), 

970 "state_type": schemas.states.StateType.SCHEDULED, 

971 "state_name": "Scheduled", 

972 "next_scheduled_start_time": date, 

973 "expected_start_time": date, 

974 "created_by": { 

975 "id": deployment_schedule.id, 

976 "display_value": deployment_schedule.slug 

977 or deployment_schedule.schedule.__class__.__name__, 

978 "type": "SCHEDULE", 

979 }, 

980 } 

981 ) 

982 

983 return runs 

984 

985 

986@db_injector 

987async def _insert_scheduled_flow_runs( 

988 db: PrefectDBInterface, session: AsyncSession, runs: list[dict[str, Any]] 

989) -> Sequence[UUID]: 

990 """ 

991 Given a list of flow runs to schedule, as generated by `_generate_scheduled_flow_runs`, 

992 inserts them into the database. Note this is a separate method to facilitate batch 

993 operations on many scheduled runs. 

994 

995 Args: 

996 session: a database session 

997 runs: a list of dicts representing flow runs to insert 

998 

999 Returns: 

1000 a list of flow run ids that were created 

1001 """ 

1002 

1003 if not runs: 1003 ↛ 1009line 1003 didn't jump to line 1009 because the condition on line 1003 was always true

1004 return [] 

1005 

1006 # gracefully insert the flow runs against the idempotency key 

1007 # this syntax (insert statement, values to insert) is most efficient 

1008 # because it uses a single bind parameter 

1009 await session.execute( 

1010 db.queries.insert(db.FlowRun).on_conflict_do_nothing( 

1011 index_elements=db.orm.flow_run_unique_upsert_columns 

1012 ), 

1013 runs, 

1014 ) 

1015 

1016 # query for the rows that were newly inserted (by checking for any flow runs with 

1017 # no corresponding flow run states) 

1018 inserted_rows = sa.select(db.FlowRun.id).where( 

1019 db.FlowRun.id.in_([r["id"] for r in runs]), 

1020 ~select(db.FlowRunState.id) 

1021 .where(db.FlowRunState.flow_run_id == db.FlowRun.id) 

1022 .exists(), 

1023 ) 

1024 inserted_flow_run_ids = (await session.execute(inserted_rows)).scalars().all() 

1025 

1026 # insert flow run states that correspond to the newly-insert rows 

1027 insert_flow_run_states: list[dict[str, Any]] = [ 

1028 {"id": uuid7(), "flow_run_id": r["id"], **r["state"]} 

1029 for r in runs 

1030 if r["id"] in inserted_flow_run_ids 

1031 ] 

1032 if insert_flow_run_states: 

1033 # this syntax (insert statement, values to insert) is most efficient 

1034 # because it uses a single bind parameter 

1035 await session.execute( 

1036 db.FlowRunState.__table__.insert(), # type: ignore[attr-defined] 

1037 insert_flow_run_states, 

1038 ) 

1039 

1040 # set the `state_id` on the newly inserted runs 

1041 stmt = db.queries.set_state_id_on_inserted_flow_runs_statement( 

1042 inserted_flow_run_ids=inserted_flow_run_ids, 

1043 insert_flow_run_states=insert_flow_run_states, 

1044 ) 

1045 

1046 await session.execute(stmt) 

1047 

1048 return inserted_flow_run_ids 

1049 

1050 

1051@db_injector 

1052async def check_work_queues_for_deployment( 

1053 db: PrefectDBInterface, session: AsyncSession, deployment_id: UUID 

1054) -> Sequence[orm_models.WorkQueue]: 

1055 """ 

1056 Get work queues that can pick up the specified deployment. 

1057 

1058 Work queues will pick up a deployment when all of the following are met. 

1059 

1060 - The deployment has ALL tags that the work queue has (i.e. the work 

1061 queue's tags must be a subset of the deployment's tags). 

1062 - The work queue's specified deployment IDs match the deployment's ID, 

1063 or the work queue does NOT have specified deployment IDs. 

1064 - The work queue's specified flow runners match the deployment's flow 

1065 runner or the work queue does NOT have a specified flow runner. 

1066 

1067 Notes on the query: 

1068 

1069 - Our database currently allows either "null" and empty lists as 

1070 null values in filters, so we need to catch both cases with "or". 

1071 - `A.contains(B)` should be interpreted as "True if A 

1072 contains B". 

1073 

1074 Returns: 

1075 List[orm_models.WorkQueue]: WorkQueues 

1076 """ 

1077 deployment = await session.get(db.Deployment, deployment_id) 

1078 if not deployment: 

1079 raise ObjectNotFoundError(f"Deployment with id {deployment_id} not found") 

1080 

1081 def json_contains(a: Any, b: Any) -> sa.ColumnElement[bool]: 

1082 return sa.type_coerce(a, type_=JSONB).contains(sa.type_coerce(b, type_=JSONB)) 

1083 

1084 query = ( 

1085 select(db.WorkQueue) 

1086 # work queue tags are a subset of deployment tags 

1087 .filter( 

1088 or_( 

1089 json_contains(deployment.tags, db.WorkQueue.filter["tags"]), 

1090 json_contains([], db.WorkQueue.filter["tags"]), 

1091 json_contains(None, db.WorkQueue.filter["tags"]), 

1092 ) 

1093 ) 

1094 # deployment_ids is null or contains the deployment's ID 

1095 .filter( 

1096 or_( 

1097 json_contains( 

1098 db.WorkQueue.filter["deployment_ids"], 

1099 str(deployment.id), 

1100 ), 

1101 json_contains(None, db.WorkQueue.filter["deployment_ids"]), 

1102 json_contains([], db.WorkQueue.filter["deployment_ids"]), 

1103 ) 

1104 ) 

1105 ) 

1106 

1107 result = await session.execute(query) 

1108 return result.scalars().unique().all() 

1109 

1110 

1111@db_injector 

1112async def create_deployment_schedules( 

1113 db: PrefectDBInterface, 

1114 session: AsyncSession, 

1115 deployment_id: UUID, 

1116 schedules: list[schemas.actions.DeploymentScheduleCreate], 

1117) -> list[schemas.core.DeploymentSchedule]: 

1118 """ 

1119 Creates a deployment's schedules. 

1120 

1121 Args: 

1122 session: A database session 

1123 deployment_id: a deployment id 

1124 schedules: a list of deployment schedule create actions 

1125 """ 

1126 

1127 schedules_with_deployment_id: list[dict[str, Any]] = [] 

1128 for schedule in schedules: 

1129 # Exclude 'replaces' as it's a deploy-time directive, not a persisted field 

1130 data = schedule.model_dump(exclude={"replaces"}) 

1131 data["deployment_id"] = deployment_id 

1132 schedules_with_deployment_id.append(data) 

1133 

1134 models = [ 

1135 db.DeploymentSchedule(**schedule) for schedule in schedules_with_deployment_id 

1136 ] 

1137 session.add_all(models) 

1138 await session.flush() 

1139 

1140 return [ 

1141 schemas.core.DeploymentSchedule.model_validate(m, from_attributes=True) 

1142 for m in models 

1143 ] 

1144 

1145 

1146@db_injector 

1147async def read_deployment_schedules( 

1148 db: PrefectDBInterface, 

1149 session: AsyncSession, 

1150 deployment_id: UUID, 

1151 deployment_schedule_filter: Optional[ 

1152 schemas.filters.DeploymentScheduleFilter 

1153 ] = None, 

1154) -> list[schemas.core.DeploymentSchedule]: 

1155 """ 

1156 Reads a deployment's schedules. 

1157 

1158 Args: 

1159 session: A database session 

1160 deployment_id: a deployment id 

1161 

1162 Returns: 

1163 list[schemas.core.DeploymentSchedule]: the deployment's schedules 

1164 """ 

1165 

1166 query = ( 

1167 sa.select(db.DeploymentSchedule) 

1168 .where(db.DeploymentSchedule.deployment_id == deployment_id) 

1169 .order_by(db.DeploymentSchedule.updated.desc()) 

1170 ) 

1171 

1172 if deployment_schedule_filter: 1172 ↛ 1173line 1172 didn't jump to line 1173 because the condition on line 1172 was never true

1173 query = query.where(deployment_schedule_filter.as_sql_filter()) 

1174 

1175 result = await session.execute(query) 

1176 

1177 return [ 

1178 schemas.core.DeploymentSchedule.model_validate(s, from_attributes=True) 

1179 for s in result.scalars().all() 

1180 ] 

1181 

1182 

1183@db_injector 

1184async def update_deployment_schedule( 

1185 db: PrefectDBInterface, 

1186 session: AsyncSession, 

1187 deployment_id: UUID, 

1188 schedule: schemas.actions.DeploymentScheduleUpdate, 

1189 deployment_schedule_id: UUID | None = None, 

1190 deployment_schedule_slug: str | None = None, 

1191) -> bool: 

1192 """ 

1193 Updates a deployment's schedules. 

1194 

1195 Args: 

1196 session: A database session 

1197 deployment_schedule_id: a deployment schedule id 

1198 schedule: a deployment schedule update action 

1199 """ 

1200 # Exclude 'replaces' as it's a deploy-time directive, not a persisted field 

1201 update_values = schedule.model_dump(exclude_none=True, exclude={"replaces"}) 

1202 if deployment_schedule_id: 

1203 result = await session.execute( 

1204 sa.update(db.DeploymentSchedule) 

1205 .where( 

1206 sa.and_( 

1207 db.DeploymentSchedule.id == deployment_schedule_id, 

1208 db.DeploymentSchedule.deployment_id == deployment_id, 

1209 ) 

1210 ) 

1211 .values(**update_values) 

1212 ) 

1213 elif deployment_schedule_slug: 

1214 result = await session.execute( 

1215 sa.update(db.DeploymentSchedule) 

1216 .where( 

1217 sa.and_( 

1218 db.DeploymentSchedule.slug == deployment_schedule_slug, 

1219 db.DeploymentSchedule.deployment_id == deployment_id, 

1220 ) 

1221 ) 

1222 .values(**update_values) 

1223 ) 

1224 else: 

1225 raise ValueError( 

1226 "Either deployment_schedule_id or deployment_schedule_slug must be provided" 

1227 ) 

1228 

1229 return result.rowcount > 0 

1230 

1231 

1232@db_injector 

1233async def delete_schedules_for_deployment( 

1234 db: PrefectDBInterface, session: AsyncSession, deployment_id: UUID 

1235) -> bool: 

1236 """ 

1237 Deletes a deployment schedule. 

1238 

1239 Args: 

1240 session: A database session 

1241 deployment_id: a deployment id 

1242 """ 

1243 

1244 deployment = await session.get(db.Deployment, deployment_id) 

1245 assert deployment is not None 

1246 

1247 result = await session.execute( 

1248 sa.delete(db.DeploymentSchedule).where( 

1249 db.DeploymentSchedule.deployment_id == deployment_id 

1250 ) 

1251 ) 

1252 

1253 await session.refresh(deployment) 

1254 return result.rowcount > 0 

1255 

1256 

1257@db_injector 

1258async def delete_deployment_schedule( 

1259 db: PrefectDBInterface, 

1260 session: AsyncSession, 

1261 deployment_id: UUID, 

1262 deployment_schedule_id: UUID, 

1263) -> bool: 

1264 """ 

1265 Deletes a deployment schedule. 

1266 

1267 Args: 

1268 session: A database session 

1269 deployment_schedule_id: a deployment schedule id 

1270 """ 

1271 

1272 result = await session.execute( 

1273 sa.delete(db.DeploymentSchedule).where( 

1274 sa.and_( 

1275 db.DeploymentSchedule.id == deployment_schedule_id, 

1276 db.DeploymentSchedule.deployment_id == deployment_id, 

1277 ) 

1278 ) 

1279 ) 

1280 

1281 return result.rowcount > 0 

1282 

1283 

1284async def mark_deployments_ready( 

1285 *, 

1286 db: PrefectDBInterface = Depends(provide_database_interface), 

1287 deployment_ids: Optional[Iterable[UUID]] = None, 

1288 work_queue_ids: Optional[Iterable[UUID]] = None, 

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

1290) -> None: 

1291 deployment_ids = deployment_ids or [] 

1292 work_queue_ids = work_queue_ids or [] 

1293 

1294 if not deployment_ids and not work_queue_ids: 

1295 return 

1296 

1297 async with db.session_context( 

1298 begin_transaction=True, 

1299 ) as session: 

1300 # ORDER BY id locks rows in deterministic order so concurrent 

1301 # calls cannot deadlock. SKIP LOCKED is intentionally avoided — 

1302 # a fresh poll must wait for a concurrent stale transition, 

1303 # not no-op and let the stale transition overwrite it. 

1304 locked = ( 

1305 select(db.Deployment.id, db.Deployment.status) 

1306 .where( 

1307 sa.or_( 

1308 db.Deployment.id.in_(deployment_ids), 

1309 db.Deployment.work_queue_id.in_(work_queue_ids), 

1310 ), 

1311 ) 

1312 .order_by(db.Deployment.id) 

1313 .with_for_update() 

1314 .cte("locked") 

1315 ) 

1316 

1317 result = await session.execute(select(locked.c.id, locked.c.status)) 

1318 rows = result.all() 

1319 

1320 if not rows: 

1321 return 

1322 

1323 unready_deployments = [ 

1324 row.id for row in rows if row.status == DeploymentStatus.NOT_READY 

1325 ] 

1326 

1327 last_polled = now("UTC") 

1328 

1329 # keeps `updated` untouched to not trigger recent schedules calculation 

1330 await session.execute( 

1331 sa.update(db.Deployment) 

1332 .where(db.Deployment.id.in_(select(locked.c.id))) 

1333 .values( 

1334 status=DeploymentStatus.READY, 

1335 last_polled=last_polled, 

1336 updated=db.Deployment.updated, 

1337 ) 

1338 ) 

1339 

1340 if not unready_deployments: 

1341 return 

1342 

1343 async with PrefectServerEventsClient() as events: 

1344 for deployment_id in unready_deployments: 

1345 await events.emit( 

1346 await deployment_status_event( 

1347 session=session, 

1348 deployment_id=deployment_id, 

1349 status=DeploymentStatus.READY, 

1350 occurred=last_polled, 

1351 ) 

1352 ) 

1353 

1354 

1355@db_injector 

1356async def mark_deployments_not_ready( 

1357 db: PrefectDBInterface, 

1358 deployment_ids: Optional[Iterable[UUID]] = None, 

1359 work_queue_ids: Optional[Iterable[UUID]] = None, 

1360) -> None: 

1361 try: 

1362 deployment_ids = deployment_ids or [] 

1363 work_queue_ids = work_queue_ids or [] 

1364 

1365 if not deployment_ids and not work_queue_ids: 1365 ↛ 1368line 1365 didn't jump to line 1368 because the condition on line 1365 was always true

1366 return 

1367 

1368 async with db.session_context( 

1369 begin_transaction=True, 

1370 ) as session: 

1371 # See comment in mark_deployments_ready. 

1372 locked = ( 

1373 select(db.Deployment.id, db.Deployment.status) 

1374 .where( 

1375 sa.or_( 

1376 db.Deployment.id.in_(deployment_ids), 

1377 db.Deployment.work_queue_id.in_(work_queue_ids), 

1378 ), 

1379 ) 

1380 .order_by(db.Deployment.id) 

1381 .with_for_update() 

1382 .cte("locked") 

1383 ) 

1384 

1385 result = await session.execute(select(locked.c.id, locked.c.status)) 

1386 rows = result.all() 

1387 

1388 if not rows: 

1389 return 

1390 

1391 ready_deployments = [ 

1392 row.id for row in rows if row.status == DeploymentStatus.READY 

1393 ] 

1394 

1395 # keeps `updated` untouched to not trigger recent schedules calculation 

1396 await session.execute( 

1397 sa.update(db.Deployment) 

1398 .where(db.Deployment.id.in_(select(locked.c.id))) 

1399 .values( 

1400 status=DeploymentStatus.NOT_READY, updated=db.Deployment.updated 

1401 ) 

1402 ) 

1403 

1404 if not ready_deployments: 

1405 return 

1406 

1407 async with PrefectServerEventsClient() as events: 

1408 for deployment_id in ready_deployments: 

1409 await events.emit( 

1410 await deployment_status_event( 

1411 session=session, 

1412 deployment_id=deployment_id, 

1413 status=DeploymentStatus.NOT_READY, 

1414 occurred=now("UTC"), 

1415 ) 

1416 ) 

1417 except Exception as exc: 

1418 logger.error(f"Error marking deployments as not ready: {exc}", exc_info=True) 

1419 

1420 

1421async def with_system_labels_for_deployment( 

1422 session: AsyncSession, 

1423 deployment: schemas.core.Deployment, 

1424) -> schemas.core.KeyValueLabels: 

1425 """Augment user supplied labels with system default labels for a deployment.""" 

1426 default_labels = cast( 

1427 schemas.core.KeyValueLabels, 

1428 { 

1429 "prefect.flow.id": str(deployment.flow_id), 

1430 }, 

1431 ) 

1432 

1433 user_supplied_labels = deployment.labels or {} 

1434 

1435 parent_labels = ( 

1436 await models.flows.read_flow_labels(session, deployment.flow_id) 

1437 ) or {} 

1438 

1439 return parent_labels | default_labels | user_supplied_labels 

1440 

1441 

1442async def with_system_labels_for_deployment_flow_run( 

1443 session: AsyncSession, 

1444 deployment: orm_models.Deployment, 

1445 user_supplied_labels: Optional[schemas.core.KeyValueLabels] = None, 

1446) -> schemas.core.KeyValueLabels: 

1447 """Generate system labels for a flow run created from a deployment. 

1448 

1449 Args: 

1450 session: Database session 

1451 deployment: The deployment the flow run is created from 

1452 user_supplied_labels: Optional user-supplied labels to include 

1453 

1454 Returns: 

1455 Complete set of labels for the flow run 

1456 """ 

1457 system_labels = cast( 

1458 schemas.core.KeyValueLabels, 

1459 { 

1460 "prefect.flow.id": str(deployment.flow_id), 

1461 "prefect.deployment.id": str(deployment.id), 

1462 }, 

1463 ) 

1464 

1465 # Use deployment labels as parent labels for flow runs 

1466 parent_labels = deployment.labels or {} 

1467 user_labels = user_supplied_labels or {} 

1468 

1469 return parent_labels | system_labels | user_labels 

1470 

1471 

1472def _detect_deployment_changed_fields( 

1473 old_snapshot: dict[str, Any], 

1474 new: orm_models.Deployment, 

1475) -> dict[str, dict[str, Any]]: 

1476 """Compare a snapshot of old field values with the new deployment ORM object.""" 

1477 changed_fields: dict[str, dict[str, Any]] = {} 

1478 for field in DEPLOYMENT_EVENT_FIELDS: 

1479 old_value = old_snapshot.get(field) 

1480 new_value = getattr(new, field, None) 

1481 if old_value != new_value: 

1482 changed_fields[field] = { 

1483 "from": old_value, 

1484 "to": new_value, 

1485 } 

1486 return changed_fields 

1487 

1488 

1489async def emit_deployment_created_event( 

1490 session: AsyncSession, 

1491 deployment: orm_models.Deployment, 

1492) -> None: 

1493 """Emit an event when a deployment is created.""" 

1494 async with PrefectServerEventsClient() as events_client: 

1495 await events_client.emit( 

1496 await deployment_created_event( 

1497 session=session, 

1498 deployment=deployment, 

1499 occurred=now("UTC"), 

1500 ) 

1501 ) 

1502 

1503 

1504async def emit_deployment_updated_event( 

1505 session: AsyncSession, 

1506 deployment: orm_models.Deployment, 

1507 changed_fields: dict[str, dict[str, Any]], 

1508) -> None: 

1509 """Emit an event when a deployment is updated.""" 

1510 if not changed_fields: 

1511 return 

1512 async with PrefectServerEventsClient() as events_client: 

1513 await events_client.emit( 

1514 await deployment_updated_event( 

1515 session=session, 

1516 deployment=deployment, 

1517 changed_fields=changed_fields, 

1518 occurred=now("UTC"), 

1519 ) 

1520 ) 

1521 

1522 

1523async def emit_deployment_deleted_event( 

1524 session: AsyncSession, 

1525 deployment: orm_models.Deployment, 

1526) -> None: 

1527 """Emit an event when a deployment is deleted.""" 

1528 async with PrefectServerEventsClient() as events_client: 

1529 await events_client.emit( 

1530 await deployment_deleted_event( 

1531 session=session, 

1532 deployment=deployment, 

1533 occurred=now("UTC"), 

1534 ) 

1535 )