Coverage for /usr/local/lib/python3.12/site-packages/prefect/server/schemas/core.py: 93%
403 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"""
2Full schemas of Prefect REST API objects.
3"""
5from __future__ import annotations
7import datetime
8from typing import (
9 TYPE_CHECKING,
10 Annotated,
11 Any,
12 ClassVar,
13 Dict,
14 List,
15 Optional,
16 Type,
17 Union,
18)
19from uuid import UUID
21from pydantic import (
22 AfterValidator,
23 BaseModel,
24 ConfigDict,
25 Field,
26 StrictBool,
27 StrictFloat,
28 StrictInt,
29 field_validator,
30 model_validator,
31)
32from sqlalchemy.ext.asyncio import AsyncSession
33from typing_extensions import Literal, Self
35from prefect._internal.schemas.validators import (
36 get_or_create_run_name,
37 list_length_50_or_less,
38 set_run_policy_deprecated_fields,
39 validate_cache_key_length,
40 validate_default_queue_id_not_none,
41 validate_max_metadata_length,
42 validate_name_present_on_nonanonymous_blocks,
43 validate_not_negative,
44 validate_parent_and_ref_diff,
45 validate_schedule_max_scheduled_runs,
46)
47from prefect.server.schemas import schedules, states
48from prefect.server.schemas.statuses import WorkPoolStatus
49from prefect.server.utilities.schemas.bases import (
50 ORMBaseModel,
51 PrefectBaseModel,
52 TimeSeriesBaseModel,
53)
54from prefect.settings import PREFECT_DEPLOYMENT_SCHEDULE_MAX_SCHEDULED_RUNS
55from prefect.types import (
56 MAX_VARIABLE_NAME_LENGTH,
57 DateTime,
58 LaxUrl,
59 Name,
60 NameOrEmpty,
61 NonEmptyishName,
62 NonNegativeInteger,
63 PositiveInteger,
64 StrictVariableValue,
65)
66from prefect.types._datetime import now
67from prefect.types.names import raise_on_name_alphanumeric_dashes_only
68from prefect.utilities.collections import (
69 AutoEnum,
70 dict_to_flatdict,
71 flatdict_to_dict,
72)
73from prefect.utilities.names import generate_slug, obfuscate
75if TYPE_CHECKING: 75 ↛ 76line 75 didn't jump to line 76 because the condition on line 75 was never true
76 from prefect.server.database import orm_models
78DEFAULT_BLOCK_SCHEMA_VERSION = "non-versioned"
80KeyValueLabels = dict[str, Union[StrictBool, StrictInt, StrictFloat, str]]
83class Flow(ORMBaseModel):
84 """An ORM representation of flow data."""
86 name: Name = Field(
87 default=..., description="The name of the flow", examples=["my-flow"]
88 )
89 tags: List[str] = Field(
90 default_factory=list,
91 description="A list of flow tags",
92 examples=[["tag-1", "tag-2"]],
93 )
94 labels: Union[KeyValueLabels, None] = Field(
95 default_factory=dict,
96 description="A dictionary of key-value labels. Values can be strings, numbers, or booleans.",
97 examples=[{"key": "value1", "key2": 42}],
98 )
101class FlowRunPolicy(PrefectBaseModel):
102 """Defines of how a flow run should retry."""
104 max_retries: int = Field(
105 default=0,
106 description=(
107 "The maximum number of retries. Field is not used. Please use `retries`"
108 " instead."
109 ),
110 deprecated=True,
111 )
112 retry_delay_seconds: float = Field(
113 default=0,
114 description=(
115 "The delay between retries. Field is not used. Please use `retry_delay`"
116 " instead."
117 ),
118 deprecated=True,
119 )
120 retries: Optional[int] = Field(default=None, description="The number of retries.")
121 retry_delay: Optional[int] = Field(
122 default=None, description="The delay time between retries, in seconds."
123 )
124 pause_keys: Optional[set[str]] = Field(
125 default_factory=set, description="Tracks pauses this run has observed."
126 )
127 resuming: Optional[bool] = Field(
128 default=False, description="Indicates if this run is resuming from a pause."
129 )
130 retry_type: Optional[Literal["in_process", "reschedule"]] = Field(
131 default=None, description="The type of retry this run is undergoing."
132 )
134 @model_validator(mode="before")
135 def populate_deprecated_fields(cls, values: dict[str, Any]) -> dict[str, Any]:
136 return set_run_policy_deprecated_fields(values)
139class CreatedBy(BaseModel):
140 id: Optional[UUID] = Field(
141 default=None, description="The id of the creator of the object."
142 )
143 type: Optional[str] = Field(
144 default=None, description="The type of the creator of the object."
145 )
146 display_value: Optional[str] = Field(
147 default=None, description="The display value for the creator."
148 )
151class UpdatedBy(BaseModel):
152 id: Optional[UUID] = Field(
153 default=None, description="The id of the updater of the object."
154 )
155 type: Optional[str] = Field(
156 default=None, description="The type of the updater of the object."
157 )
158 display_value: Optional[str] = Field(
159 default=None, description="The display value for the updater."
160 )
163class ConcurrencyLimitStrategy(AutoEnum):
164 """
165 Enumeration of concurrency collision strategies.
166 """
168 ENQUEUE = AutoEnum.auto()
169 CANCEL_NEW = AutoEnum.auto()
172class ConcurrencyOptions(BaseModel):
173 """
174 Class for storing the concurrency config in database.
175 """
177 collision_strategy: ConcurrencyLimitStrategy
178 grace_period_seconds: Optional[int] = Field(
179 default=None,
180 ge=60,
181 le=86400,
182 description="Grace period in seconds for infrastructure to start before concurrency slots are revoked. If not set, falls back to server setting.",
183 )
186class FlowRun(TimeSeriesBaseModel, ORMBaseModel):
187 """An ORM representation of flow run data."""
189 name: str = Field(
190 default_factory=lambda: generate_slug(2),
191 description=(
192 "The name of the flow run. Defaults to a random slug if not specified."
193 ),
194 examples=["my-flow-run"],
195 )
196 flow_id: UUID = Field(default=..., description="The id of the flow being run.")
197 state_id: Optional[UUID] = Field(
198 default=None, description="The id of the flow run's current state."
199 )
200 deployment_id: Optional[UUID] = Field(
201 default=None,
202 description=(
203 "The id of the deployment associated with this flow run, if available."
204 ),
205 )
206 deployment_version: Optional[str] = Field(
207 default=None,
208 description="The version of the deployment associated with this flow run.",
209 examples=["1.0"],
210 )
211 work_queue_name: Optional[str] = Field(
212 default=None, description="The work queue that handled this flow run."
213 )
214 flow_version: Optional[str] = Field(
215 default=None,
216 description="The version of the flow executed in this flow run.",
217 examples=["1.0"],
218 )
219 parameters: Dict[str, Any] = Field(
220 default_factory=dict, description="Parameters for the flow run."
221 )
222 idempotency_key: Optional[str] = Field(
223 default=None,
224 description=(
225 "An optional idempotency key for the flow run. Used to ensure the same flow"
226 " run is not created multiple times."
227 ),
228 )
229 context: Dict[str, Any] = Field(
230 default_factory=dict,
231 description="Additional context for the flow run.",
232 examples=[{"my_var": "my_value"}],
233 )
234 empirical_policy: FlowRunPolicy = Field(
235 default_factory=FlowRunPolicy,
236 )
237 tags: List[str] = Field(
238 default_factory=list,
239 description="A list of tags on the flow run",
240 examples=[["tag-1", "tag-2"]],
241 )
242 labels: Union[KeyValueLabels, None] = Field(
243 default_factory=dict,
244 description="A dictionary of key-value labels. Values can be strings, numbers, or booleans.",
245 examples=[{"key": "value1", "key2": 42}],
246 )
247 parent_task_run_id: Optional[UUID] = Field(
248 default=None,
249 description=(
250 "If the flow run is a subflow, the id of the 'dummy' task in the parent"
251 " flow used to track subflow state."
252 ),
253 )
255 state_type: Optional[states.StateType] = Field(
256 default=None, description="The type of the current flow run state."
257 )
258 state_name: Optional[str] = Field(
259 default=None, description="The name of the current flow run state."
260 )
261 run_count: int = Field(
262 default=0, description="The number of times the flow run was executed."
263 )
264 expected_start_time: Optional[DateTime] = Field(
265 default=None,
266 description="The flow run's expected start time.",
267 )
268 next_scheduled_start_time: Optional[DateTime] = Field(
269 default=None,
270 description="The next time the flow run is scheduled to start.",
271 )
272 start_time: Optional[DateTime] = Field(
273 default=None, description="The actual start time."
274 )
275 end_time: Optional[DateTime] = Field(
276 default=None, description="The actual end time."
277 )
278 total_run_time: datetime.timedelta = Field(
279 default=datetime.timedelta(0),
280 description=(
281 "Total run time. If the flow run was executed multiple times, the time of"
282 " each run will be summed."
283 ),
284 )
285 estimated_run_time: datetime.timedelta = Field(
286 default=datetime.timedelta(0),
287 description="A real-time estimate of the total run time.",
288 )
289 estimated_start_time_delta: datetime.timedelta = Field(
290 default=datetime.timedelta(0),
291 description="The difference between actual and expected start time.",
292 )
293 auto_scheduled: bool = Field(
294 default=False,
295 description="Whether or not the flow run was automatically scheduled.",
296 )
297 infrastructure_document_id: Optional[UUID] = Field(
298 default=None,
299 description="The block document defining infrastructure to use this flow run.",
300 )
301 infrastructure_pid: Optional[str] = Field(
302 default=None,
303 description="The id of the flow run as returned by an infrastructure block.",
304 )
305 created_by: Optional[CreatedBy] = Field(
306 default=None,
307 description="Optional information about the creator of this flow run.",
308 )
309 work_queue_id: Optional[UUID] = Field(
310 default=None, description="The id of the run's work pool queue."
311 )
313 # relationships
314 # flow: Flow = None
315 # task_runs: List["TaskRun"] = Field(default_factory=list)
316 state: Optional[states.State] = Field(
317 default=None, description="The current state of the flow run."
318 )
319 # parent_task_run: "TaskRun" = None
321 job_variables: Optional[Dict[str, Any]] = Field(
322 default=None,
323 description="Variables used as overrides in the base job template",
324 )
326 @field_validator("name", mode="before")
327 @classmethod
328 def set_name(cls, name: str) -> str:
329 return get_or_create_run_name(name)
331 def __eq__(self, other: Any) -> bool:
332 """
333 Check for "equality" to another flow run schema
335 Estimates times are rolling and will always change with repeated queries for
336 a flow run so we ignore them during equality checks.
337 """
338 if isinstance(other, FlowRun):
339 exclude_fields = {"estimated_run_time", "estimated_start_time_delta"}
340 return self.model_dump(exclude=exclude_fields) == other.model_dump(
341 exclude=exclude_fields
342 )
343 return super().__eq__(other)
346class TaskRunPolicy(PrefectBaseModel):
347 """Defines of how a task run should retry."""
349 max_retries: int = Field(
350 default=0,
351 description=(
352 "The maximum number of retries. Field is not used. Please use `retries`"
353 " instead."
354 ),
355 deprecated=True,
356 )
357 retry_delay_seconds: float = Field(
358 default=0,
359 description=(
360 "The delay between retries. Field is not used. Please use `retry_delay`"
361 " instead."
362 ),
363 deprecated=True,
364 )
365 retries: Optional[int] = Field(default=None, description="The number of retries.")
366 retry_delay: Union[None, int, float, List[int], List[float]] = Field(
367 default=None,
368 description="A delay time or list of delay times between retries, in seconds.",
369 )
370 retry_jitter_factor: Optional[float] = Field(
371 default=None, description="Determines the amount a retry should jitter"
372 )
374 @model_validator(mode="before")
375 def populate_deprecated_fields(cls, values: dict[str, Any]) -> dict[str, Any]:
376 return set_run_policy_deprecated_fields(values)
378 @field_validator("retry_delay")
379 @classmethod
380 def validate_configured_retry_delays(
381 cls, v: int | float | list[int] | list[float] | None
382 ) -> int | float | list[int] | list[float] | None:
383 return list_length_50_or_less(v)
385 @field_validator("retry_jitter_factor")
386 @classmethod
387 def validate_jitter_factor(cls, v: float | None) -> float | None:
388 return validate_not_negative(v)
391class RunInput(PrefectBaseModel):
392 """
393 Base class for classes that represent inputs to runs, which
394 could include, constants, parameters, task runs or flow runs.
395 """
397 model_config: ClassVar[ConfigDict] = ConfigDict(frozen=True)
399 input_type: str
402class TaskRunResult(RunInput):
403 """Represents a task run result input to another task run."""
405 input_type: Literal["task_run"] = "task_run"
406 id: UUID
409class FlowRunResult(RunInput):
410 input_type: Literal["flow_run"] = "flow_run"
411 id: UUID
414class Parameter(RunInput):
415 """Represents a parameter input to a task run."""
417 input_type: Literal["parameter"] = "parameter"
418 name: str
421class Constant(RunInput):
422 """Represents constant input value to a task run."""
424 input_type: Literal["constant"] = "constant"
425 type: str
428class TaskRun(TimeSeriesBaseModel, ORMBaseModel):
429 """An ORM representation of task run data."""
431 name: str = Field(
432 default_factory=lambda: generate_slug(2), examples=["my-task-run"]
433 )
434 flow_run_id: Optional[UUID] = Field(
435 default=None, description="The flow run id of the task run."
436 )
437 task_key: str = Field(
438 default=..., description="A unique identifier for the task being run."
439 )
440 dynamic_key: str = Field(
441 default=...,
442 description=(
443 "A dynamic key used to differentiate between multiple runs of the same task"
444 " within the same flow run."
445 ),
446 )
447 cache_key: Optional[str] = Field(
448 default=None,
449 description=(
450 "An optional cache key. If a COMPLETED state associated with this cache key"
451 " is found, the cached COMPLETED state will be used instead of executing"
452 " the task run."
453 ),
454 )
455 cache_expiration: Optional[DateTime] = Field(
456 default=None, description="Specifies when the cached state should expire."
457 )
458 task_version: Optional[str] = Field(
459 default=None, description="The version of the task being run."
460 )
461 empirical_policy: TaskRunPolicy = Field(
462 default_factory=TaskRunPolicy,
463 )
464 tags: List[str] = Field(
465 default_factory=list,
466 description="A list of tags for the task run.",
467 examples=[["tag-1", "tag-2"]],
468 )
469 labels: Union[KeyValueLabels, None] = Field(
470 default_factory=dict,
471 description="A dictionary of key-value labels. Values can be strings, numbers, or booleans.",
472 examples=[{"key": "value1", "key2": 42}],
473 )
474 state_id: Optional[UUID] = Field(
475 default=None, description="The id of the current task run state."
476 )
477 task_inputs: Dict[
478 str, List[Union[TaskRunResult, FlowRunResult, Parameter, Constant]]
479 ] = Field(
480 default_factory=dict,
481 description=(
482 "Tracks the source of inputs to a task run. Used for internal bookkeeping."
483 ),
484 )
485 state_type: Optional[states.StateType] = Field(
486 default=None, description="The type of the current task run state."
487 )
488 state_name: Optional[str] = Field(
489 default=None, description="The name of the current task run state."
490 )
491 run_count: int = Field(
492 default=0, description="The number of times the task run has been executed."
493 )
494 flow_run_run_count: int = Field(
495 default=0,
496 description=(
497 "If the parent flow has retried, this indicates the flow retry this run is"
498 " associated with."
499 ),
500 )
501 expected_start_time: Optional[DateTime] = Field(
502 default=None,
503 description="The task run's expected start time.",
504 )
506 # the next scheduled start time will be populated
507 # whenever the run is in a scheduled state
508 next_scheduled_start_time: Optional[DateTime] = Field(
509 default=None,
510 description="The next time the task run is scheduled to start.",
511 )
512 start_time: Optional[DateTime] = Field(
513 default=None, description="The actual start time."
514 )
515 end_time: Optional[DateTime] = Field(
516 default=None, description="The actual end time."
517 )
518 total_run_time: datetime.timedelta = Field(
519 default=datetime.timedelta(0),
520 description=(
521 "Total run time. If the task run was executed multiple times, the time of"
522 " each run will be summed."
523 ),
524 )
525 estimated_run_time: datetime.timedelta = Field(
526 default=datetime.timedelta(0),
527 description="A real-time estimate of total run time.",
528 )
529 estimated_start_time_delta: datetime.timedelta = Field(
530 default=datetime.timedelta(0),
531 description="The difference between actual and expected start time.",
532 )
534 # relationships
535 # flow_run: FlowRun = None
536 # subflow_runs: List[FlowRun] = Field(default_factory=list)
537 state: Optional[states.State] = Field(
538 default=None, description="The current task run state."
539 )
541 @field_validator("name", mode="before")
542 @classmethod
543 def set_name(cls, name: str) -> str:
544 return get_or_create_run_name(name)
546 @field_validator("cache_key")
547 @classmethod
548 def validate_cache_key(cls, cache_key: str) -> str:
549 return validate_cache_key_length(cache_key)
552class DeploymentSchedule(ORMBaseModel):
553 deployment_id: Optional[UUID] = Field(
554 default=None,
555 description="The deployment id associated with this schedule.",
556 )
557 schedule: schedules.SCHEDULE_TYPES = Field(
558 default=..., description="The schedule for the deployment."
559 )
560 active: bool = Field(
561 default=True, description="Whether or not the schedule is active."
562 )
563 max_scheduled_runs: Optional[PositiveInteger] = Field(
564 default=None,
565 description="The maximum number of scheduled runs for the schedule.",
566 )
567 parameters: dict[str, Any] = Field(
568 default_factory=dict, description="A dictionary of parameter value overrides."
569 )
570 slug: Optional[str] = Field(
571 default=None,
572 description="A unique slug for the schedule.",
573 )
575 @field_validator("max_scheduled_runs")
576 @classmethod
577 def validate_max_scheduled_runs(cls, v: int) -> int:
578 return validate_schedule_max_scheduled_runs(
579 v, PREFECT_DEPLOYMENT_SCHEDULE_MAX_SCHEDULED_RUNS.value()
580 )
583class VersionInfo(PrefectBaseModel, extra="allow"):
584 type: str = Field(default=..., description="The type of version info.")
585 version: str = Field(default=..., description="The version of the deployment.")
588class Deployment(ORMBaseModel):
589 """An ORM representation of deployment data."""
591 model_config: ClassVar[ConfigDict] = ConfigDict(populate_by_name=True)
593 name: NameOrEmpty = Field(default=..., description="The name of the deployment.")
594 version: Optional[str] = Field(
595 default=None, description="An optional version for the deployment."
596 )
597 description: Optional[str] = Field(
598 default=None, description="A description for the deployment."
599 )
600 flow_id: UUID = Field(
601 default=..., description="The flow id associated with the deployment."
602 )
603 paused: bool = Field(
604 default=False, description="Whether or not the deployment is paused."
605 )
606 schedules: list[DeploymentSchedule] = Field(
607 default_factory=lambda: [],
608 description="A list of schedules for the deployment.",
609 )
610 concurrency_limit: Optional[NonNegativeInteger] = Field(
611 default=None, description="The concurrency limit for the deployment."
612 )
613 concurrency_limit_id: Optional[UUID] = Field(
614 default=None,
615 description="The concurrency limit id associated with the deployment.",
616 )
617 concurrency_options: Optional[ConcurrencyOptions] = Field(
618 default=None, description="The concurrency options for the deployment."
619 )
620 job_variables: Dict[str, Any] = Field(
621 default_factory=dict,
622 description="Overrides to apply to flow run infrastructure at runtime.",
623 )
624 parameters: Dict[str, Any] = Field(
625 default_factory=dict,
626 description="Parameters for flow runs scheduled by the deployment.",
627 )
628 pull_steps: Optional[list[dict[str, Any]]] = Field(
629 default=None,
630 description="Pull steps for cloning and running this deployment.",
631 )
632 tags: List[str] = Field(
633 default_factory=list,
634 description="A list of tags for the deployment",
635 examples=[["tag-1", "tag-2"]],
636 )
637 labels: Union[KeyValueLabels, None] = Field(
638 default_factory=dict,
639 description="A dictionary of key-value labels. Values can be strings, numbers, or booleans.",
640 examples=[{"key": "value1", "key2": 42}],
641 )
642 work_queue_name: Optional[str] = Field(
643 default=None,
644 description=(
645 "The work queue for the deployment. If no work queue is set, work will not"
646 " be scheduled."
647 ),
648 )
649 last_polled: Optional[DateTime] = Field(
650 default=None,
651 description="The last time the deployment was polled for status updates.",
652 )
653 parameter_openapi_schema: Optional[Dict[str, Any]] = Field(
654 default_factory=dict,
655 description="The parameter schema of the flow, including defaults.",
656 )
657 path: Optional[str] = Field(
658 default=None,
659 description=(
660 "The path to the working directory for the workflow, relative to remote"
661 " storage or an absolute path."
662 ),
663 )
664 entrypoint: Optional[str] = Field(
665 default=None,
666 description=(
667 "The path to the entrypoint for the workflow, relative to the `path`."
668 ),
669 )
670 storage_document_id: Optional[UUID] = Field(
671 default=None,
672 description="The block document defining storage used for this flow.",
673 )
674 infrastructure_document_id: Optional[UUID] = Field(
675 default=None,
676 description="The block document defining infrastructure to use for flow runs.",
677 )
678 created_by: Optional[CreatedBy] = Field(
679 default=None,
680 description="Optional information about the creator of this deployment.",
681 )
682 updated_by: Optional[UpdatedBy] = Field(
683 default=None,
684 description="Optional information about the updater of this deployment.",
685 )
686 work_queue_id: Optional[UUID] = Field(
687 default=None,
688 description=(
689 "The id of the work pool queue to which this deployment is assigned."
690 ),
691 )
692 enforce_parameter_schema: bool = Field(
693 default=True,
694 description=(
695 "Whether or not the deployment should enforce the parameter schema."
696 ),
697 )
700class ConcurrencyLimit(ORMBaseModel):
701 """An ORM representation of a concurrency limit."""
703 tag: str = Field(
704 default=..., description="A tag the concurrency limit is applied to."
705 )
706 concurrency_limit: int = Field(default=..., description="The concurrency limit.")
707 active_slots: list[UUID] = Field(
708 default_factory=lambda: [],
709 description="A list of active run ids using a concurrency slot",
710 )
713class ConcurrencyLimitV2(ORMBaseModel):
714 """An ORM representation of a v2 concurrency limit."""
716 active: bool = Field(
717 default=True, description="Whether the concurrency limit is active."
718 )
719 name: Name = Field(default=..., description="The name of the concurrency limit.")
720 limit: int = Field(default=..., description="The concurrency limit.")
721 active_slots: int = Field(default=0, description="The number of active slots.")
722 denied_slots: int = Field(default=0, description="The number of denied slots.")
723 slot_decay_per_second: float = Field(
724 default=0,
725 description="The decay rate for active slots when used as a rate limit.",
726 )
727 avg_slot_occupancy_seconds: float = Field(
728 default=2.0, description="The average amount of time a slot is occupied."
729 )
732class BlockType(ORMBaseModel):
733 """An ORM representation of a block type"""
735 name: Name = Field(default=..., description="A block type's name")
736 slug: str = Field(default=..., description="A block type's slug")
737 logo_url: Optional[LaxUrl] = Field( # TODO: make it HttpUrl
738 default=None, description="Web URL for the block type's logo"
739 )
740 documentation_url: Optional[LaxUrl] = Field( # TODO: make it HttpUrl
741 default=None, description="Web URL for the block type's documentation"
742 )
743 description: Optional[str] = Field(
744 default=None,
745 description="A short blurb about the corresponding block's intended use",
746 )
747 code_example: Optional[str] = Field(
748 default=None,
749 description="A code snippet demonstrating use of the corresponding block",
750 )
751 is_protected: bool = Field(
752 default=False, description="Protected block types cannot be modified via API."
753 )
756class BlockSchema(ORMBaseModel):
757 """An ORM representation of a block schema."""
759 checksum: str = Field(default=..., description="The block schema's unique checksum")
760 fields: Dict[str, Any] = Field(
761 default_factory=dict,
762 description="The block schema's field schema",
763 json_schema_extra={"additionalProperties": True},
764 )
765 block_type_id: Optional[UUID] = Field(default=..., description="A block type ID")
766 block_type: Optional[BlockType] = Field(
767 default=None, description="The associated block type"
768 )
769 capabilities: List[str] = Field(
770 default_factory=list,
771 description="A list of Block capabilities",
772 )
773 version: str = Field(
774 default=DEFAULT_BLOCK_SCHEMA_VERSION,
775 description="Human readable identifier for the block schema",
776 )
779class BlockSchemaReference(ORMBaseModel):
780 """An ORM representation of a block schema reference."""
782 parent_block_schema_id: UUID = Field(
783 default=..., description="ID of block schema the reference is nested within"
784 )
785 parent_block_schema: Optional[BlockSchema] = Field(
786 default=None, description="The block schema the reference is nested within"
787 )
788 reference_block_schema_id: UUID = Field(
789 default=..., description="ID of the nested block schema"
790 )
791 reference_block_schema: Optional[BlockSchema] = Field(
792 default=None, description="The nested block schema"
793 )
794 name: str = Field(
795 default=..., description="The name that the reference is nested under"
796 )
799class BlockDocument(ORMBaseModel):
800 """An ORM representation of a block document."""
802 name: Optional[Name] = Field(
803 default=None,
804 description=(
805 "The block document's name. Not required for anonymous block documents."
806 ),
807 )
808 data: dict[str, Any] = Field(
809 default_factory=dict, description="The block document's data"
810 )
811 block_schema_id: UUID = Field(default=..., description="A block schema ID")
812 block_schema: Optional[BlockSchema] = Field(
813 default=None, description="The associated block schema"
814 )
815 block_type_id: UUID = Field(default=..., description="A block type ID")
816 block_type_name: Optional[str] = Field(
817 default=None, description="The associated block type's name"
818 )
819 block_type: Optional[BlockType] = Field(
820 default=None, description="The associated block type"
821 )
822 block_document_references: dict[str, dict[str, Any]] = Field(
823 default_factory=dict, description="Record of the block document's references"
824 )
825 is_anonymous: bool = Field(
826 default=False,
827 description=(
828 "Whether the block is anonymous (anonymous blocks are usually created by"
829 " Prefect automatically)"
830 ),
831 )
833 @model_validator(mode="before")
834 def validate_name_is_present_if_not_anonymous(
835 cls, values: dict[str, Any]
836 ) -> dict[str, Any]:
837 return validate_name_present_on_nonanonymous_blocks(values)
839 @classmethod
840 async def from_orm_model(
841 cls: type[Self],
842 session: AsyncSession,
843 orm_block_document: "orm_models.ORMBlockDocument",
844 include_secrets: bool = False,
845 ) -> Self:
846 data = await orm_block_document.decrypt_data(session=session)
847 # if secrets are not included, obfuscate them based on the schema's
848 # `secret_fields`. Note this walks any nested blocks as well. If the
849 # nested blocks were recovered from named blocks, they will already
850 # be obfuscated, but if nested fields were hardcoded into the parent
851 # blocks data, this is the only opportunity to obfuscate them.
852 if not include_secrets:
853 flat_data = dict_to_flatdict(data)
854 # iterate over the (possibly nested) secret fields
855 # and obfuscate their data
856 for secret_field in orm_block_document.block_schema.fields.get(
857 "secret_fields", []
858 ):
859 secret_key = tuple(secret_field.split("."))
860 if flat_data.get(secret_key) is not None: 860 ↛ 861line 860 didn't jump to line 861 because the condition on line 860 was never true
861 flat_data[secret_key] = obfuscate(flat_data[secret_key])
862 # If a wildcard (*) is in the current secret key path, we take the portion
863 # of the path before the wildcard and compare it to the same level of each
864 # key. A match means that the field is nested under the secret key and should
865 # be obfuscated.
866 elif "*" in secret_key: 866 ↛ 867line 866 didn't jump to line 867 because the condition on line 866 was never true
867 wildcard_index = secret_key.index("*")
868 for data_key in flat_data.keys():
869 if secret_key[0:wildcard_index] == data_key[0:wildcard_index]:
870 flat_data[data_key] = obfuscate(flat_data[data_key])
871 data = flatdict_to_dict(flat_data)
872 return cls(
873 id=orm_block_document.id,
874 created=orm_block_document.created,
875 updated=orm_block_document.updated,
876 name=orm_block_document.name,
877 data=data,
878 block_schema_id=orm_block_document.block_schema_id,
879 block_schema=orm_block_document.block_schema,
880 block_type_id=orm_block_document.block_type_id,
881 block_type_name=orm_block_document.block_type_name,
882 block_type=orm_block_document.block_type,
883 is_anonymous=orm_block_document.is_anonymous,
884 )
887class BlockDocumentReference(ORMBaseModel):
888 """An ORM representation of a block document reference."""
890 parent_block_document_id: UUID = Field(
891 default=..., description="ID of block document the reference is nested within"
892 )
893 parent_block_document: Optional[BlockDocument] = Field(
894 default=None, description="The block document the reference is nested within"
895 )
896 reference_block_document_id: UUID = Field(
897 default=..., description="ID of the nested block document"
898 )
899 reference_block_document: Optional[BlockDocument] = Field(
900 default=None, description="The nested block document"
901 )
902 name: str = Field(
903 default=..., description="The name that the reference is nested under"
904 )
906 @model_validator(mode="before")
907 def validate_parent_and_ref_are_different(
908 cls, values: dict[str, Any]
909 ) -> dict[str, Any]:
910 return validate_parent_and_ref_diff(values)
913class Configuration(ORMBaseModel):
914 """An ORM representation of account info."""
916 key: str = Field(default=..., description="Account info key")
917 value: Dict[str, Any] = Field(default=..., description="Account info")
920class ServerDefaultResultStorage(PrefectBaseModel):
921 """Server-side default result storage configuration."""
923 default_result_storage_block_id: Optional[UUID] = Field(
924 default=None,
925 description="The block document ID of the server default result storage block.",
926 )
929class ServerDefaultResultStorageUpdate(PrefectBaseModel):
930 """Request payload for setting the server default result storage block."""
932 default_result_storage_block_id: UUID = Field(
933 default=...,
934 description="The block document ID of the server default result storage block.",
935 )
938class SavedSearchFilter(PrefectBaseModel):
939 """A filter for a saved search model. Intended for use by the Prefect UI."""
941 object: str = Field(default=..., description="The object over which to filter.")
942 property: str = Field(
943 default=..., description="The property of the object on which to filter."
944 )
945 type: str = Field(default=..., description="The type of the property.")
946 operation: str = Field(
947 default=...,
948 description="The operator to apply to the object. For example, `equals`.",
949 )
950 value: Any = Field(
951 default=..., description="A JSON-compatible value for the filter."
952 )
955class SavedSearch(ORMBaseModel):
956 """An ORM representation of saved search data. Represents a set of filter criteria."""
958 name: str = Field(default=..., description="The name of the saved search.")
959 filters: list[SavedSearchFilter] = Field(
960 default_factory=lambda: [],
961 description="The filter set for the saved search.",
962 )
965class Log(TimeSeriesBaseModel, ORMBaseModel):
966 """An ORM representation of log data."""
968 name: str = Field(default=..., description="The logger name.")
969 level: int = Field(default=..., description="The log level.")
970 message: str = Field(default=..., description="The log message.")
971 timestamp: DateTime = Field(default=..., description="The log timestamp.")
972 flow_run_id: Optional[UUID] = Field(
973 default=None, description="The flow run ID associated with the log."
974 )
975 task_run_id: Optional[UUID] = Field(
976 default=None, description="The task run ID associated with the log."
977 )
980class QueueFilter(PrefectBaseModel):
981 """Filter criteria definition for a work queue."""
983 tags: Optional[list[str]] = Field(
984 default=None,
985 description="Only include flow runs with these tags in the work queue.",
986 )
987 deployment_ids: Optional[list[UUID]] = Field(
988 default=None,
989 description="Only include flow runs from these deployments in the work queue.",
990 )
993class WorkQueue(ORMBaseModel):
994 """An ORM representation of a work queue"""
996 name: Name = Field(default=..., description="The name of the work queue.")
997 description: Optional[str] = Field(
998 default="", description="An optional description for the work queue."
999 )
1000 is_paused: bool = Field(
1001 default=False, description="Whether or not the work queue is paused."
1002 )
1003 concurrency_limit: Optional[NonNegativeInteger] = Field(
1004 default=None, description="An optional concurrency limit for the work queue."
1005 )
1006 priority: PositiveInteger = Field(
1007 default=1,
1008 description=(
1009 "The queue's priority. Lower values are higher priority (1 is the highest)."
1010 ),
1011 )
1012 # Will be required after a future migration
1013 work_pool_id: Optional[UUID] = Field(
1014 default=None, description="The work pool with which the queue is associated."
1015 )
1016 filter: Optional[QueueFilter] = Field(
1017 default=None,
1018 description="DEPRECATED: Filter criteria for the work queue.",
1019 deprecated=True,
1020 )
1021 last_polled: Optional[DateTime] = Field(
1022 default=None, description="The last time an agent polled this queue for work."
1023 )
1026class WorkQueueHealthPolicy(PrefectBaseModel):
1027 maximum_late_runs: Optional[int] = Field(
1028 default=0,
1029 description=(
1030 "The maximum number of late runs in the work queue before it is deemed"
1031 " unhealthy. Defaults to `0`."
1032 ),
1033 )
1034 maximum_seconds_since_last_polled: Optional[int] = Field(
1035 default=60,
1036 description=(
1037 "The maximum number of time in seconds elapsed since work queue has been"
1038 " polled before it is deemed unhealthy. Defaults to `60`."
1039 ),
1040 )
1042 def evaluate_health_status(
1043 self, late_runs_count: int, last_polled: Optional[DateTime] = None
1044 ) -> bool:
1045 """
1046 Given empirical information about the state of the work queue, evaluate its health status.
1048 Args:
1049 late_runs_count: the count of late runs for the work queue.
1050 last_polled: the last time the work queue was polled, if available.
1052 Returns:
1053 bool: whether or not the work queue is healthy.
1054 """
1055 healthy = True
1056 if ( 1056 ↛ 1060line 1056 didn't jump to line 1060 because the condition on line 1056 was never true
1057 self.maximum_late_runs is not None
1058 and late_runs_count > self.maximum_late_runs
1059 ):
1060 healthy = False
1062 if self.maximum_seconds_since_last_polled is not None: 1062 ↛ 1070line 1062 didn't jump to line 1070 because the condition on line 1062 was always true
1063 if (
1064 last_polled is None
1065 or (now("UTC") - last_polled).total_seconds()
1066 > self.maximum_seconds_since_last_polled
1067 ):
1068 healthy = False
1070 return healthy
1073class WorkQueueStatusDetail(PrefectBaseModel):
1074 healthy: bool = Field(..., description="Whether or not the work queue is healthy.")
1075 late_runs_count: int = Field(
1076 default=0, description="The number of late flow runs in the work queue."
1077 )
1078 last_polled: Optional[DateTime] = Field(
1079 default=None, description="The last time an agent polled this queue for work."
1080 )
1081 health_check_policy: WorkQueueHealthPolicy = Field(
1082 ...,
1083 description=(
1084 "The policy used to determine whether or not the work queue is healthy."
1085 ),
1086 )
1089class Agent(ORMBaseModel):
1090 """An ORM representation of an agent"""
1092 name: str = Field(
1093 default_factory=lambda: generate_slug(2),
1094 description=(
1095 "The name of the agent. If a name is not provided, it will be"
1096 " auto-generated."
1097 ),
1098 )
1099 work_queue_id: UUID = Field(
1100 default=..., description="The work queue with which the agent is associated."
1101 )
1102 last_activity_time: Optional[DateTime] = Field(
1103 default=None, description="The last time this agent polled for work."
1104 )
1107class WorkPoolStorageConfiguration(PrefectBaseModel):
1108 """A representation of a work pool's storage configuration"""
1110 model_config: ClassVar[ConfigDict] = ConfigDict(extra="forbid")
1112 bundle_upload_step: Optional[dict[str, Any]] = Field(
1113 default=None,
1114 description="The step to use for uploading bundles to storage.",
1115 )
1116 bundle_execution_step: Optional[dict[str, Any]] = Field(
1117 default=None,
1118 description="The step to use for executing bundles.",
1119 )
1120 default_result_storage_block_id: Optional[UUID] = Field(
1121 default=None,
1122 description="The block document ID of the default result storage block.",
1123 )
1126class WorkPool(ORMBaseModel):
1127 """An ORM representation of a work pool"""
1129 name: NonEmptyishName = Field(
1130 description="The name of the work pool.",
1131 )
1132 description: Optional[str] = Field(
1133 default=None, description="A description of the work pool."
1134 )
1135 type: str = Field(description="The work pool type.")
1136 base_job_template: Dict[str, Any] = Field(
1137 default_factory=dict, description="The work pool's base job template."
1138 )
1139 is_paused: bool = Field(
1140 default=False,
1141 description="Pausing the work pool stops the delivery of all work.",
1142 )
1143 concurrency_limit: Optional[NonNegativeInteger] = Field(
1144 default=None, description="A concurrency limit for the work pool."
1145 )
1146 status: Optional[WorkPoolStatus] = Field(
1147 default=None, description="The current status of the work pool."
1148 )
1150 # this required field has a default of None so that the custom validator
1151 # below will be called and produce a more helpful error message
1152 default_queue_id: Optional[UUID] = Field(
1153 default=None, description="The id of the pool's default queue."
1154 )
1156 storage_configuration: WorkPoolStorageConfiguration = Field(
1157 default_factory=WorkPoolStorageConfiguration,
1158 description="The storage configuration for the work pool.",
1159 )
1161 @field_validator("default_queue_id")
1162 def helpful_error_for_missing_default_queue_id(cls, v: UUID | None) -> UUID:
1163 return validate_default_queue_id_not_none(v)
1165 @classmethod
1166 def model_validate(
1167 cls: Type[Self],
1168 obj: Any,
1169 *,
1170 strict: Optional[bool] = None,
1171 from_attributes: Optional[bool] = None,
1172 context: Optional[dict[str, Any]] = None,
1173 ) -> Self:
1174 parsed: WorkPool = super().model_validate(
1175 obj, strict=strict, from_attributes=from_attributes, context=context
1176 )
1177 if from_attributes: 1177 ↛ 1180line 1177 didn't jump to line 1180 because the condition on line 1177 was always true
1178 if obj.type == "prefect-agent":
1179 parsed.status = None
1180 return parsed
1183class Worker(ORMBaseModel):
1184 """An ORM representation of a worker"""
1186 name: str = Field(description="The name of the worker.")
1187 work_pool_id: UUID = Field(
1188 description="The work pool with which the queue is associated."
1189 )
1190 last_heartbeat_time: Optional[datetime.datetime] = Field(
1191 None, description="The last time the worker process sent a heartbeat."
1192 )
1193 heartbeat_interval_seconds: Optional[int] = Field(
1194 default=None,
1195 description=(
1196 "The number of seconds to expect between heartbeats sent by the worker."
1197 ),
1198 )
1201Flow.model_rebuild()
1202FlowRun.model_rebuild()
1205class Artifact(ORMBaseModel):
1206 key: Optional[str] = Field(
1207 default=None, description="An optional unique reference key for this artifact."
1208 )
1209 type: Optional[str] = Field(
1210 default=None,
1211 description=(
1212 "An identifier that describes the shape of the data field. e.g. 'result',"
1213 " 'table', 'markdown'"
1214 ),
1215 )
1216 description: Optional[str] = Field(
1217 default=None, description="A markdown-enabled description of the artifact."
1218 )
1219 # data will eventually be typed as `Optional[Union[Result, Any]]`
1220 data: Optional[Union[Dict[str, Any], Any]] = Field(
1221 default=None,
1222 description=(
1223 "Data associated with the artifact, e.g. a result.; structure depends on"
1224 " the artifact type."
1225 ),
1226 )
1227 metadata_: Optional[dict[str, str]] = Field(
1228 default=None,
1229 description=(
1230 "User-defined artifact metadata. Content must be string key and value"
1231 " pairs."
1232 ),
1233 )
1234 flow_run_id: Optional[UUID] = Field(
1235 default=None, description="The flow run associated with the artifact."
1236 )
1237 task_run_id: Optional[UUID] = Field(
1238 default=None, description="The task run associated with the artifact."
1239 )
1241 @classmethod
1242 def from_result(cls, data: Any | dict[str, Any]) -> "Artifact":
1243 artifact_info: dict[str, Any] = dict()
1244 if isinstance(data, dict):
1245 artifact_key = data.pop("artifact_key", None)
1246 if artifact_key: 1246 ↛ 1247line 1246 didn't jump to line 1247 because the condition on line 1246 was never true
1247 artifact_info["key"] = artifact_key
1249 artifact_type = data.pop("artifact_type", None)
1250 if artifact_type: 1250 ↛ 1251line 1250 didn't jump to line 1251 because the condition on line 1250 was never true
1251 artifact_info["type"] = artifact_type
1253 description = data.pop("artifact_description", None)
1254 if description: 1254 ↛ 1255line 1254 didn't jump to line 1255 because the condition on line 1254 was never true
1255 artifact_info["description"] = description
1257 return cls(data=data, **artifact_info)
1259 @field_validator("metadata_")
1260 @classmethod
1261 def validate_metadata_length(cls, v: dict[str, str]) -> dict[str, str]:
1262 return validate_max_metadata_length(v)
1265class ArtifactCollection(ORMBaseModel):
1266 key: str = Field(description="An optional unique reference key for this artifact.")
1267 latest_id: UUID = Field(
1268 description="The latest artifact ID associated with the key."
1269 )
1270 type: Optional[str] = Field(
1271 default=None,
1272 description=(
1273 "An identifier that describes the shape of the data field. e.g. 'result',"
1274 " 'table', 'markdown'"
1275 ),
1276 )
1277 description: Optional[str] = Field(
1278 default=None, description="A markdown-enabled description of the artifact."
1279 )
1280 data: Optional[Union[Dict[str, Any], Any]] = Field(
1281 default=None,
1282 description=(
1283 "Data associated with the artifact, e.g. a result.; structure depends on"
1284 " the artifact type."
1285 ),
1286 )
1287 metadata_: Optional[Dict[str, str]] = Field(
1288 default=None,
1289 description=(
1290 "User-defined artifact metadata. Content must be string key and value"
1291 " pairs."
1292 ),
1293 )
1294 flow_run_id: Optional[UUID] = Field(
1295 default=None, description="The flow run associated with the artifact."
1296 )
1297 task_run_id: Optional[UUID] = Field(
1298 default=None, description="The task run associated with the artifact."
1299 )
1302class Variable(ORMBaseModel):
1303 name: str = Field(
1304 default=...,
1305 description="The name of the variable",
1306 examples=["my-variable"],
1307 max_length=MAX_VARIABLE_NAME_LENGTH,
1308 )
1309 value: StrictVariableValue = Field(
1310 default=...,
1311 description="The value of the variable",
1312 examples=["my-value"],
1313 )
1314 tags: List[str] = Field(
1315 default_factory=list,
1316 description="A list of variable tags",
1317 examples=[["tag-1", "tag-2"]],
1318 )
1321class FlowRunInput(ORMBaseModel):
1322 flow_run_id: UUID = Field(description="The flow run ID associated with the input.")
1323 key: Annotated[str, AfterValidator(raise_on_name_alphanumeric_dashes_only)] = Field(
1324 description="The key of the input."
1325 )
1326 value: str = Field(description="The value of the input.")
1327 sender: Optional[str] = Field(default=None, description="The sender of the input.")
1330class CsrfToken(ORMBaseModel):
1331 token: str = Field(
1332 default=...,
1333 description="The CSRF token",
1334 )
1335 client: str = Field(
1336 default=..., description="The client id associated with the CSRF token"
1337 )
1338 expiration: DateTime = Field(
1339 default=..., description="The expiration time of the CSRF token"
1340 )