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
« 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"""
6from __future__ import annotations
8import datetime
9import logging
10from collections.abc import Iterable, Sequence
11from typing import TYPE_CHECKING, Any, Optional, TypeVar, cast
12from uuid import UUID
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
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
47T = TypeVar("T", bound=tuple[Any, ...])
49logger: logging.Logger = get_logger("prefect.server.models.deployments")
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}
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.
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 )
96 if auto_scheduled_only:
97 delete_query = delete_query.where(
98 db.FlowRun.auto_scheduled.is_(True),
99 )
101 if future_only:
102 delete_query = delete_query.where(
103 db.FlowRun.next_scheduled_start_time > now("UTC"),
104 )
106 await session.execute(delete_query)
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.
117 Args:
118 session: a database session
119 deployment: a deployment model
121 Returns:
122 orm_models.Deployment: the newly-created or updated deployment
124 """
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")
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)
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]
155 deployment.labels = await with_system_labels_for_deployment(session, deployment)
157 schedules = deployment.schedules
158 insert_values = deployment.model_dump_for_orm(
159 exclude_unset=True, exclude={"schedules", "version_info"}
160 )
162 requested_concurrency_limit = insert_values.pop("concurrency_limit", "unset")
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
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
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 )
194 await session.execute(insert_stmt)
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()
207 if not deployment_id:
208 return None
210 # Because this was possibly an upsert, we need to delete any existing
211 # schedules and any runs from the old deployment.
213 await _delete_scheduled_runs(
214 session=session,
215 deployment_id=deployment_id,
216 auto_scheduled_only=True,
217 future_only=True,
218 )
220 await delete_schedules_for_deployment(session=session, deployment_id=deployment_id)
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 )
237 if requested_concurrency_limit != "unset":
238 await _create_or_update_deployment_concurrency_limit(
239 db, session, deployment_id, deployment.concurrency_limit
240 )
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()
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 )
272 return result_deployment
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.
284 Args:
285 session: a database session
286 deployment_id: the ID of the deployment to modify
287 deployment: changes to a deployment model
289 Returns:
290 bool: whether the deployment was updated
292 """
294 from prefect.server.api.workers import WorkerLookups
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)
309 schedules = deployment.schedules
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 )
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")
323 if requested_global_concurrency_limit_update != "unset":
324 update_data["concurrency_limit_id"] = requested_global_concurrency_limit_update
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
332 should_update_schedules = update_data.pop("schedules", None) is not None
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
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)
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 )
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 )
405 if requested_concurrency_limit_update != "unset":
406 await _create_or_update_deployment_concurrency_limit(
407 db, session, deployment_id, deployment.concurrency_limit
408 )
410 updated = result.rowcount > 0
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 )
427 return updated
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
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
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
458 session.add(deployment)
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.
467 Args:
468 session: A database session
469 deployment_id: a deployment id
471 Returns:
472 orm_models.Deployment: the deployment
473 """
475 return await session.get(db.Deployment, deployment_id)
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.
484 Args:
485 session: A database session
486 name: a deployment name
487 flow_name: the name of the flow the deployment belongs to
489 Returns:
490 orm_models.Deployment: the deployment
491 """
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()
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 """
521 if deployment_filter:
522 query = query.where(deployment_filter.as_sql_filter())
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 )
530 query = query.where(flow_exists_clause.exists())
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 )
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())
547 query = query.where(flow_run_exists_clause.exists())
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 )
554 if work_queue_filter:
555 work_pool_exists_clause = work_pool_exists_clause.where(
556 work_queue_filter.as_sql_filter()
557 )
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())
565 query = query.where(work_pool_exists_clause.exists())
567 return query
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.
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.
599 Returns:
600 list[orm_models.Deployment]: deployments
601 """
603 query = select(db.Deployment).order_by(*sort.as_sql_sort())
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 )
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)
621 result = await session.execute(query)
622 return result.scalars().unique().all()
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.
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
648 Returns:
649 int: the number of deployments matching filters
650 """
652 query = select(sa.func.count(None)).select_from(db.Deployment)
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 )
665 result = await session.execute(query)
666 return result.scalar_one()
669@db_injector
670async def delete_deployment(
671 db: PrefectDBInterface, session: AsyncSession, deployment_id: UUID
672) -> bool:
673 """
674 Delete a deployment by id.
676 Args:
677 session: A database session
678 deployment_id: a deployment id
680 Returns:
681 bool: whether or not the deployment was deleted
682 """
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 )
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 )
697 await _delete_related_concurrency_limit(
698 db, session=session, deployment_id=deployment_id
699 )
701 result = await session.execute(
702 delete(db.Deployment).where(db.Deployment.id == deployment_id)
703 )
704 deleted = result.rowcount > 0
706 if deleted and delete_event is not None:
707 async with PrefectServerEventsClient() as events_client:
708 await events_client.emit(delete_event)
710 return deleted
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 )
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.
735 Args:
736 session: A database session
737 deployment_ids: a list of deployment ids to delete
739 Returns:
740 List[UUID]: the IDs of the deployments that were deleted
741 """
742 if not deployment_ids:
743 return []
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())
751 if not existing_deployments:
752 return []
754 existing_ids = [d.id for d in existing_deployments]
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 )
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 )
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 )
781 # Delete all deployments
782 await session.execute(
783 delete(db.Deployment).where(db.Deployment.id.in_(existing_ids))
784 )
786 # Emit delete events
787 async with PrefectServerEventsClient() as events_client:
788 for event in delete_events:
789 await events_client.emit(event)
791 return existing_ids
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
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
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.
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
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
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
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)
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.
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.
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
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.
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
905 Returns:
906 a list of dictionary representations of the `FlowRun` objects to schedule
907 """
908 runs: list[dict[str, Any]] = []
910 deployment = await session.get(db.Deployment, deployment_id)
912 if not deployment:
913 return []
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 )
923 for deployment_schedule in active_deployment_schedules:
924 dates: list[DateTime] = []
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)
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
936 tags = deployment.tags
937 if auto_scheduled:
938 tags = ["auto-scheduled"] + tags
940 parameters = {
941 **deployment.parameters,
942 **deployment_schedule.parameters,
943 }
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 )
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 )
983 return runs
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.
995 Args:
996 session: a database session
997 runs: a list of dicts representing flow runs to insert
999 Returns:
1000 a list of flow run ids that were created
1001 """
1003 if not runs: 1003 ↛ 1009line 1003 didn't jump to line 1009 because the condition on line 1003 was always true
1004 return []
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 )
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()
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 )
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 )
1046 await session.execute(stmt)
1048 return inserted_flow_run_ids
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.
1058 Work queues will pick up a deployment when all of the following are met.
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.
1067 Notes on the query:
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".
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")
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))
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 )
1107 result = await session.execute(query)
1108 return result.scalars().unique().all()
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.
1121 Args:
1122 session: A database session
1123 deployment_id: a deployment id
1124 schedules: a list of deployment schedule create actions
1125 """
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)
1134 models = [
1135 db.DeploymentSchedule(**schedule) for schedule in schedules_with_deployment_id
1136 ]
1137 session.add_all(models)
1138 await session.flush()
1140 return [
1141 schemas.core.DeploymentSchedule.model_validate(m, from_attributes=True)
1142 for m in models
1143 ]
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.
1158 Args:
1159 session: A database session
1160 deployment_id: a deployment id
1162 Returns:
1163 list[schemas.core.DeploymentSchedule]: the deployment's schedules
1164 """
1166 query = (
1167 sa.select(db.DeploymentSchedule)
1168 .where(db.DeploymentSchedule.deployment_id == deployment_id)
1169 .order_by(db.DeploymentSchedule.updated.desc())
1170 )
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())
1175 result = await session.execute(query)
1177 return [
1178 schemas.core.DeploymentSchedule.model_validate(s, from_attributes=True)
1179 for s in result.scalars().all()
1180 ]
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.
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 )
1229 return result.rowcount > 0
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.
1239 Args:
1240 session: A database session
1241 deployment_id: a deployment id
1242 """
1244 deployment = await session.get(db.Deployment, deployment_id)
1245 assert deployment is not None
1247 result = await session.execute(
1248 sa.delete(db.DeploymentSchedule).where(
1249 db.DeploymentSchedule.deployment_id == deployment_id
1250 )
1251 )
1253 await session.refresh(deployment)
1254 return result.rowcount > 0
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.
1267 Args:
1268 session: A database session
1269 deployment_schedule_id: a deployment schedule id
1270 """
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 )
1281 return result.rowcount > 0
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 []
1294 if not deployment_ids and not work_queue_ids:
1295 return
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 )
1317 result = await session.execute(select(locked.c.id, locked.c.status))
1318 rows = result.all()
1320 if not rows:
1321 return
1323 unready_deployments = [
1324 row.id for row in rows if row.status == DeploymentStatus.NOT_READY
1325 ]
1327 last_polled = now("UTC")
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 )
1340 if not unready_deployments:
1341 return
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 )
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 []
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
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 )
1385 result = await session.execute(select(locked.c.id, locked.c.status))
1386 rows = result.all()
1388 if not rows:
1389 return
1391 ready_deployments = [
1392 row.id for row in rows if row.status == DeploymentStatus.READY
1393 ]
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 )
1404 if not ready_deployments:
1405 return
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)
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 )
1433 user_supplied_labels = deployment.labels or {}
1435 parent_labels = (
1436 await models.flows.read_flow_labels(session, deployment.flow_id)
1437 ) or {}
1439 return parent_labels | default_labels | user_supplied_labels
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.
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
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 )
1465 # Use deployment labels as parent labels for flow runs
1466 parent_labels = deployment.labels or {}
1467 user_labels = user_supplied_labels or {}
1469 return parent_labels | system_labels | user_labels
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
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 )
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 )
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 )