Coverage for /usr/local/lib/python3.12/site-packages/prefect/server/schemas/actions.py: 87%
380 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"""
2Reduced schemas for accepting API actions.
3"""
5from __future__ import annotations
7import json
8from copy import deepcopy
9from typing import Annotated, Any, ClassVar, Dict, List, Optional, Union
10from uuid import UUID, uuid4
12from pydantic import (
13 AfterValidator,
14 ConfigDict,
15 Field,
16 field_validator,
17 model_validator,
18)
20import prefect.server.schemas as schemas
21from prefect._internal.schema import ParameterSchema
22from prefect._internal.schemas.validators import (
23 get_or_create_run_name,
24 normalize_schedule_rrule,
25 remove_old_deployment_fields,
26 validate_cache_key_length,
27 validate_max_metadata_length,
28 validate_name_present_on_nonanonymous_blocks,
29 validate_parameter_openapi_schema,
30 validate_parameter_size_field,
31 validate_parameters_conform_to_schema,
32 validate_parent_and_ref_diff,
33 validate_schedule_max_scheduled_runs,
34)
35from prefect.server.utilities.schemas import get_class_fields_only
36from prefect.server.utilities.schemas.bases import PrefectBaseModel
37from prefect.settings import PREFECT_DEPLOYMENT_SCHEDULE_MAX_SCHEDULED_RUNS
38from prefect.types import (
39 DateTime,
40 KeyValueLabels,
41 Name,
42 NonEmptyishName,
43 NonNegativeFloat,
44 NonNegativeInteger,
45 PositiveInteger,
46 StrictVariableValue,
47)
48from prefect.types._datetime import now
49from prefect.types.names import (
50 ArtifactKey,
51 BlockDocumentName,
52 BlockTypeSlug,
53 VariableName,
54)
55from prefect.utilities.names import generate_slug
56from prefect.utilities.templating import find_placeholders
58SizedParameters = Annotated[
59 Dict[str, Any], AfterValidator(validate_parameter_size_field)
60]
63class ActionBaseModel(PrefectBaseModel):
64 model_config: ClassVar[ConfigDict] = ConfigDict(extra="forbid")
67class FlowCreate(ActionBaseModel):
68 """Data used by the Prefect REST API to create a flow."""
70 name: Name = Field(
71 default=..., description="The name of the flow", examples=["my-flow"]
72 )
73 tags: List[str] = Field(
74 default_factory=list,
75 description="A list of flow tags",
76 examples=[["tag-1", "tag-2"]],
77 )
78 labels: Union[KeyValueLabels, None] = Field(
79 default_factory=dict,
80 description="A dictionary of key-value labels. Values can be strings, numbers, or booleans.",
81 examples=[{"key": "value1", "key2": 42}],
82 )
85class FlowUpdate(ActionBaseModel):
86 """Data used by the Prefect REST API to update a flow."""
88 tags: List[str] = Field(
89 default_factory=list,
90 description="A list of flow tags",
91 examples=[["tag-1", "tag-2"]],
92 )
95# Bare RRule schedules arriving via the API write path get an explicit
96# DTSTART injected here so the scheduler doesn't fall back to the legacy
97# 2020 anchor on every loop. The validator is attached to the *field*
98# (via Annotated) rather than to RRuleSchedule itself — if it lived on
99# the schedule class, it would also fire on every DB read and re-phase
100# INTERVAL>1 schedules. See PrefectHQ/prefect#21362.
101#
102# Fields that accept `None` (e.g. `DeploymentScheduleUpdate.schedule`)
103# use `Optional[NormalizedSchedule]` directly — Pydantic only runs the
104# `AfterValidator` on the non-None branch.
105NormalizedSchedule = Annotated[
106 schemas.schedules.SCHEDULE_TYPES, AfterValidator(normalize_schedule_rrule)
107]
110class DeploymentScheduleCreate(ActionBaseModel):
111 active: bool = Field(
112 default=True, description="Whether or not the schedule is active."
113 )
114 schedule: NormalizedSchedule = Field(
115 default=..., description="The schedule for the deployment."
116 )
117 max_scheduled_runs: Optional[PositiveInteger] = Field(
118 default=None,
119 description="The maximum number of scheduled runs for the schedule.",
120 )
121 parameters: dict[str, Any] = Field(
122 default_factory=dict, description="A dictionary of parameter value overrides."
123 )
124 slug: Optional[str] = Field(
125 default=None,
126 description="A unique identifier for the schedule.",
127 )
128 replaces: Optional[str] = Field(
129 default=None,
130 description="The slug of an existing schedule that this schedule replaces. Used for renaming slugs.",
131 )
133 @field_validator("max_scheduled_runs")
134 @classmethod
135 def validate_max_scheduled_runs(
136 cls, v: PositiveInteger | None
137 ) -> PositiveInteger | None:
138 return validate_schedule_max_scheduled_runs(
139 v, PREFECT_DEPLOYMENT_SCHEDULE_MAX_SCHEDULED_RUNS.value()
140 )
143class DeploymentScheduleUpdate(ActionBaseModel):
144 active: Optional[bool] = Field(
145 default=None, description="Whether or not the schedule is active."
146 )
147 schedule: Optional[NormalizedSchedule] = Field(
148 default=None, description="The schedule for the deployment."
149 )
151 max_scheduled_runs: Optional[PositiveInteger] = Field(
152 default=None,
153 description="The maximum number of scheduled runs for the schedule.",
154 )
155 parameters: dict[str, Any] = Field(
156 default_factory=dict, description="A dictionary of parameter value overrides."
157 )
158 slug: Optional[str] = Field(
159 default=None,
160 description="A unique identifier for the schedule.",
161 )
162 replaces: Optional[str] = Field(
163 default=None,
164 description="The slug of an existing schedule that this schedule replaces. Used for renaming slugs.",
165 )
167 @field_validator("max_scheduled_runs")
168 @classmethod
169 def validate_max_scheduled_runs(
170 cls, v: PositiveInteger | None
171 ) -> PositiveInteger | None:
172 return validate_schedule_max_scheduled_runs(
173 v, PREFECT_DEPLOYMENT_SCHEDULE_MAX_SCHEDULED_RUNS.value()
174 )
177class DeploymentCreate(ActionBaseModel):
178 """Data used by the Prefect REST API to create a deployment."""
180 name: str = Field(
181 default=...,
182 description="The name of the deployment.",
183 examples=["my-deployment"],
184 )
185 flow_id: UUID = Field(
186 default=..., description="The ID of the flow associated with the deployment."
187 )
188 paused: bool = Field(
189 default=False, description="Whether or not the deployment is paused."
190 )
191 schedules: list[DeploymentScheduleCreate] = Field(
192 default_factory=lambda: [],
193 description="A list of schedules for the deployment.",
194 )
195 concurrency_limit: Optional[PositiveInteger] = Field(
196 default=None, description="The deployment's concurrency limit."
197 )
198 concurrency_options: Optional[schemas.core.ConcurrencyOptions] = Field(
199 default=None, description="The deployment's concurrency options."
200 )
201 global_concurrency_limit_id: Optional[UUID] = Field(
202 default=None,
203 description="The ID of the global concurrency limit to apply to the deployment.",
204 )
205 enforce_parameter_schema: bool = Field(
206 default=True,
207 description=(
208 "Whether or not the deployment should enforce the parameter schema."
209 ),
210 )
211 parameter_openapi_schema: Optional[ParameterSchema] = Field(
212 default_factory=lambda: {"type": "object", "properties": {}},
213 description="The parameter schema of the flow, including defaults.",
214 json_schema_extra={"additionalProperties": True},
215 )
216 parameters: SizedParameters = Field(
217 default_factory=dict,
218 description="Parameters for flow runs scheduled by the deployment.",
219 json_schema_extra={"additionalProperties": True},
220 )
221 tags: List[str] = Field(
222 default_factory=list,
223 description="A list of deployment tags.",
224 examples=[["tag-1", "tag-2"]],
225 )
226 labels: Union[KeyValueLabels, None] = Field(
227 default_factory=dict,
228 description="A dictionary of key-value labels. Values can be strings, numbers, or booleans.",
229 examples=[{"key": "value1", "key2": 42}],
230 )
231 pull_steps: Optional[List[dict[str, Any]]] = Field(None)
233 work_queue_name: Optional[str] = Field(None)
234 work_pool_name: Optional[str] = Field(
235 default=None,
236 description="The name of the deployment's work pool.",
237 examples=["my-work-pool"],
238 )
239 storage_document_id: Optional[UUID] = Field(None)
240 infrastructure_document_id: Optional[UUID] = Field(None)
241 description: Optional[str] = Field(None)
242 path: Optional[str] = Field(None)
243 version: Optional[str] = Field(None)
244 entrypoint: Optional[str] = Field(None)
245 job_variables: Dict[str, Any] = Field(
246 default_factory=dict,
247 description="Overrides for the flow's infrastructure configuration.",
248 json_schema_extra={"additionalProperties": True},
249 )
251 version_info: Optional[schemas.core.VersionInfo] = Field(
252 default=None, description="A description of this version of the deployment."
253 )
255 def check_valid_configuration(self, base_job_template: dict[str, Any]) -> None:
256 """
257 Check that the combination of base_job_template defaults and job_variables
258 conforms to the specified schema.
260 NOTE: This method does not hydrate block references in default values within the
261 base job template to validate them. Failing to do this can cause user-facing
262 errors. Instead of this method, use `validate_job_variables_for_deployment`
263 function from `prefect_cloud.orion.api.validation`.
264 """
265 # This import is here to avoid a circular import
266 from prefect.utilities.schema_tools import validate
268 variables_schema = deepcopy(base_job_template.get("variables"))
270 if variables_schema is not None:
271 validate(
272 self.job_variables,
273 variables_schema,
274 raise_on_error=True,
275 preprocess=True,
276 ignore_required=True,
277 )
279 @model_validator(mode="before")
280 @classmethod
281 def remove_old_fields(cls, values: dict[str, Any]) -> dict[str, Any]:
282 return remove_old_deployment_fields(values)
284 @model_validator(mode="before")
285 def _validate_parameters_conform_to_schema(
286 cls, values: dict[str, Any]
287 ) -> dict[str, Any]:
288 values["parameters"] = validate_parameters_conform_to_schema(
289 values.get("parameters", {}), values
290 )
291 schema = validate_parameter_openapi_schema(
292 values.get("parameter_openapi_schema"), values
293 )
294 if schema is not None:
295 values["parameter_openapi_schema"] = schema
296 return values
298 @model_validator(mode="before")
299 def _validate_concurrency_limits(cls, values: dict[str, Any]) -> dict[str, Any]:
300 """Validate that a deployment does not have both a concurrency limit and global concurrency limit."""
301 if values.get("concurrency_limit") and values.get(
302 "global_concurrency_limit_id"
303 ):
304 raise ValueError(
305 "A deployment cannot have both a concurrency limit and a global concurrency limit."
306 )
307 return values
310class DeploymentUpdate(ActionBaseModel):
311 """Data used by the Prefect REST API to update a deployment."""
313 @model_validator(mode="before")
314 @classmethod
315 def remove_old_fields(cls, values: dict[str, Any]) -> dict[str, Any]:
316 return remove_old_deployment_fields(values)
318 version: Optional[str] = Field(default=None)
319 description: Optional[str] = Field(default=None)
320 paused: bool = Field(
321 default=False, description="Whether or not the deployment is paused."
322 )
323 schedules: list[DeploymentScheduleUpdate] = Field(
324 default_factory=lambda: [],
325 description="A list of schedules for the deployment.",
326 )
327 concurrency_limit: Optional[PositiveInteger] = Field(
328 default=None, description="The deployment's concurrency limit."
329 )
330 concurrency_options: Optional[schemas.core.ConcurrencyOptions] = Field(
331 default=None, description="The deployment's concurrency options."
332 )
333 global_concurrency_limit_id: Optional[UUID] = Field(
334 default=None,
335 description="The ID of the global concurrency limit to apply to the deployment.",
336 )
337 parameters: Optional[SizedParameters] = Field(
338 default=None,
339 description="Parameters for flow runs scheduled by the deployment.",
340 )
341 parameter_openapi_schema: Optional[ParameterSchema] = Field(
342 default=None,
343 description="The parameter schema of the flow, including defaults.",
344 )
345 tags: List[str] = Field(
346 default_factory=list,
347 description="A list of deployment tags.",
348 examples=[["tag-1", "tag-2"]],
349 )
350 work_queue_name: Optional[str] = Field(default=None)
351 work_pool_name: Optional[str] = Field(
352 default=None,
353 description="The name of the deployment's work pool.",
354 examples=["my-work-pool"],
355 )
356 path: Optional[str] = Field(default=None)
357 job_variables: Optional[Dict[str, Any]] = Field(
358 default=None,
359 description="Overrides for the flow's infrastructure configuration.",
360 )
361 pull_steps: Optional[List[dict[str, Any]]] = Field(default=None)
362 entrypoint: Optional[str] = Field(default=None)
363 storage_document_id: Optional[UUID] = Field(default=None)
364 infrastructure_document_id: Optional[UUID] = Field(default=None)
365 enforce_parameter_schema: Optional[bool] = Field(
366 default=None,
367 description=(
368 "Whether or not the deployment should enforce the parameter schema."
369 ),
370 )
371 model_config: ClassVar[ConfigDict] = ConfigDict(populate_by_name=True)
373 version_info: Optional[schemas.core.VersionInfo] = Field(
374 default=None, description="A description of this version of the deployment."
375 )
377 def check_valid_configuration(self, base_job_template: dict[str, Any]) -> None:
378 """
379 Check that the combination of base_job_template defaults and job_variables
380 conforms to the schema specified in the base_job_template.
382 NOTE: This method does not hydrate block references in default values within the
383 base job template to validate them. Failing to do this can cause user-facing
384 errors. Instead of this method, use `validate_job_variables_for_deployment`
385 function from `prefect_cloud.orion.api.validation`.
386 """
387 # This import is here to avoid a circular import
388 from prefect.utilities.schema_tools import validate
390 variables_schema = deepcopy(base_job_template.get("variables"))
392 if variables_schema is not None and self.job_variables is not None:
393 errors = validate(
394 self.job_variables,
395 variables_schema,
396 raise_on_error=False,
397 preprocess=True,
398 ignore_required=True,
399 )
400 if errors:
401 for error in errors:
402 raise error
404 @model_validator(mode="before")
405 def _validate_concurrency_limits(cls, values: dict[str, Any]) -> dict[str, Any]:
406 """Validate that a deployment does not have both a concurrency limit and global concurrency limit."""
407 if values.get("concurrency_limit") and values.get(
408 "global_concurrency_limit_id"
409 ):
410 raise ValueError(
411 "A deployment cannot have both a concurrency limit and a global concurrency limit."
412 )
413 return values
416class FlowRunUpdate(ActionBaseModel):
417 """Data used by the Prefect REST API to update a flow run."""
419 name: Optional[str] = Field(None)
420 flow_version: Optional[str] = Field(None)
421 parameters: SizedParameters = Field(default_factory=dict)
422 empirical_policy: schemas.core.FlowRunPolicy = Field(
423 default_factory=schemas.core.FlowRunPolicy
424 )
425 tags: List[str] = Field(default_factory=list)
426 infrastructure_pid: Optional[str] = Field(None)
427 job_variables: Optional[Dict[str, Any]] = Field(None)
429 @field_validator("name", mode="before")
430 @classmethod
431 def set_name(cls, name: str) -> str:
432 return get_or_create_run_name(name)
435class StateCreate(ActionBaseModel):
436 """Data used by the Prefect REST API to create a new state."""
438 type: schemas.states.StateType = Field(
439 default=..., description="The type of the state to create"
440 )
441 name: Optional[str] = Field(
442 default=None, description="The name of the state to create"
443 )
444 message: Optional[str] = Field(
445 default=None, description="The message of the state to create"
446 )
447 data: Optional[Any] = Field(
448 default=None, description="The data of the state to create"
449 )
450 state_details: schemas.states.StateDetails = Field(
451 default_factory=schemas.states.StateDetails,
452 description="The details of the state to create",
453 )
455 @model_validator(mode="after")
456 def default_name_from_type(self):
457 """If a name is not provided, use the type"""
458 # if `type` is not in `values` it means the `type` didn't pass its own
459 # validation check and an error will be raised after this function is called
460 name = self.name
461 if name is None and self.type:
462 self.name = " ".join([v.capitalize() for v in self.type.value.split("_")])
463 return self
465 @model_validator(mode="after")
466 def default_scheduled_start_time(self):
467 from prefect.server.schemas.states import StateType
469 if self.type == StateType.SCHEDULED:
470 if not self.state_details.scheduled_time:
471 self.state_details.scheduled_time = now("UTC")
473 return self
476class TaskRunCreate(ActionBaseModel):
477 """Data used by the Prefect REST API to create a task run"""
479 id: Optional[UUID] = Field(
480 default=None,
481 description="The ID to assign to the task run. If not provided, a random UUID will be generated.",
482 )
483 # TaskRunCreate states must be provided as StateCreate objects
484 state: Optional[StateCreate] = Field(
485 default=None, description="The state of the task run to create"
486 )
488 name: str = Field(
489 default_factory=lambda: generate_slug(2), examples=["my-task-run"]
490 )
491 flow_run_id: Optional[UUID] = Field(
492 default=None, description="The flow run id of the task run."
493 )
494 task_key: str = Field(
495 default=..., description="A unique identifier for the task being run."
496 )
497 dynamic_key: str = Field(
498 default=...,
499 description=(
500 "A dynamic key used to differentiate between multiple runs of the same task"
501 " within the same flow run."
502 ),
503 )
504 cache_key: Optional[str] = Field(
505 default=None,
506 description=(
507 "An optional cache key. If a COMPLETED state associated with this cache key"
508 " is found, the cached COMPLETED state will be used instead of executing"
509 " the task run."
510 ),
511 )
512 cache_expiration: Optional[DateTime] = Field(
513 default=None, description="Specifies when the cached state should expire."
514 )
515 task_version: Optional[str] = Field(
516 default=None, description="The version of the task being run."
517 )
518 empirical_policy: schemas.core.TaskRunPolicy = Field(
519 default_factory=schemas.core.TaskRunPolicy,
520 )
521 tags: List[str] = Field(
522 default_factory=list,
523 description="A list of tags for the task run.",
524 examples=[["tag-1", "tag-2"]],
525 )
526 labels: Union[KeyValueLabels, None] = Field(
527 default_factory=dict,
528 description="A dictionary of key-value labels. Values can be strings, numbers, or booleans.",
529 examples=[{"key": "value1", "key2": 42}],
530 )
531 task_inputs: Dict[
532 str,
533 List[
534 Union[
535 schemas.core.TaskRunResult,
536 schemas.core.FlowRunResult,
537 schemas.core.Parameter,
538 schemas.core.Constant,
539 ]
540 ],
541 ] = Field(
542 default_factory=dict,
543 description="The inputs to the task run.",
544 )
546 @field_validator("name", mode="before")
547 @classmethod
548 def set_name(cls, name: str) -> str:
549 return get_or_create_run_name(name)
551 @field_validator("cache_key")
552 @classmethod
553 def validate_cache_key(cls, cache_key: str | None) -> str | None:
554 return validate_cache_key_length(cache_key)
557class TaskRunUpdate(ActionBaseModel):
558 """Data used by the Prefect REST API to update a task run"""
560 name: str = Field(
561 default_factory=lambda: generate_slug(2), examples=["my-task-run"]
562 )
564 @field_validator("name", mode="before")
565 @classmethod
566 def set_name(cls, name: str) -> str:
567 return get_or_create_run_name(name)
570class FlowRunCreate(ActionBaseModel):
571 """Data used by the Prefect REST API to create a flow run."""
573 # FlowRunCreate states must be provided as StateCreate objects
574 state: Optional[StateCreate] = Field(
575 default=None, description="The state of the flow run to create"
576 )
578 name: str = Field(
579 default_factory=lambda: generate_slug(2),
580 description=(
581 "The name of the flow run. Defaults to a random slug if not specified."
582 ),
583 examples=["my-flow-run"],
584 )
585 flow_id: UUID = Field(default=..., description="The id of the flow being run.")
586 flow_version: Optional[str] = Field(
587 default=None, description="The version of the flow being run."
588 )
589 parameters: SizedParameters = Field(
590 default_factory=dict,
591 )
592 context: Dict[str, Any] = Field(
593 default_factory=dict,
594 description="The context of the flow run.",
595 )
596 parent_task_run_id: Optional[UUID] = Field(None)
597 infrastructure_document_id: Optional[UUID] = Field(None)
598 empirical_policy: schemas.core.FlowRunPolicy = Field(
599 default_factory=schemas.core.FlowRunPolicy,
600 description="The empirical policy for the flow run.",
601 )
602 tags: List[str] = Field(
603 default_factory=list,
604 description="A list of tags for the flow run.",
605 examples=[["tag-1", "tag-2"]],
606 )
607 labels: Union[KeyValueLabels, None] = Field(
608 default_factory=dict,
609 description="A dictionary of key-value labels. Values can be strings, numbers, or booleans.",
610 examples=[{"key": "value1", "key2": 42}],
611 )
612 idempotency_key: Optional[str] = Field(
613 None,
614 description=(
615 "An optional idempotency key. If a flow run with the same idempotency key"
616 " has already been created, the existing flow run will be returned."
617 ),
618 )
619 work_pool_name: Optional[str] = Field(
620 default=None,
621 description="The name of the work pool to run the flow run in.",
622 )
623 work_queue_name: Optional[str] = Field(
624 default=None,
625 description="The name of the work queue to place the flow run in.",
626 )
627 job_variables: Optional[Dict[str, Any]] = Field(
628 default=None,
629 description="The job variables to use when setting up flow run infrastructure.",
630 )
632 # DEPRECATED
634 deployment_id: Optional[UUID] = Field(
635 None,
636 description=(
637 "DEPRECATED: The id of the deployment associated with this flow run, if"
638 " available."
639 ),
640 deprecated=True,
641 )
643 @field_validator("name", mode="before")
644 @classmethod
645 def set_name(cls, name: str) -> str:
646 return get_or_create_run_name(name)
649class DeploymentFlowRunCreate(ActionBaseModel):
650 """Data used by the Prefect REST API to create a flow run from a deployment."""
652 # FlowRunCreate states must be provided as StateCreate objects
653 state: Optional[StateCreate] = Field(
654 default=None, description="The state of the flow run to create"
655 )
657 name: str = Field(
658 default_factory=lambda: generate_slug(2),
659 description=(
660 "The name of the flow run. Defaults to a random slug if not specified."
661 ),
662 examples=["my-flow-run"],
663 )
664 parameters: SizedParameters = Field(
665 default_factory=dict,
666 json_schema_extra={"additionalProperties": True},
667 )
668 enforce_parameter_schema: Optional[bool] = Field(
669 default=None,
670 description="Whether or not to enforce the parameter schema on this run.",
671 )
672 context: Dict[str, Any] = Field(default_factory=dict)
673 infrastructure_document_id: Optional[UUID] = Field(None)
674 empirical_policy: schemas.core.FlowRunPolicy = Field(
675 default_factory=schemas.core.FlowRunPolicy,
676 description="The empirical policy for the flow run.",
677 )
678 tags: List[str] = Field(
679 default_factory=list,
680 description="A list of tags for the flow run.",
681 examples=[["tag-1", "tag-2"]],
682 )
683 idempotency_key: Optional[str] = Field(
684 None,
685 description=(
686 "An optional idempotency key. If a flow run with the same idempotency key"
687 " has already been created, the existing flow run will be returned."
688 ),
689 )
690 labels: Union[KeyValueLabels, None] = Field(
691 None,
692 description="A dictionary of key-value labels. Values can be strings, numbers, or booleans.",
693 examples=[{"key": "value1", "key2": 42}],
694 )
695 parent_task_run_id: Optional[UUID] = Field(None)
696 work_queue_name: Optional[str] = Field(None)
697 job_variables: Optional[Dict[str, Any]] = Field(
698 default_factory=dict,
699 json_schema_extra={"additionalProperties": True},
700 )
702 @field_validator("name", mode="before")
703 @classmethod
704 def set_name(cls, name: str) -> str:
705 return get_or_create_run_name(name)
708class SavedSearchCreate(ActionBaseModel):
709 """Data used by the Prefect REST API to create a saved search."""
711 name: str = Field(default=..., description="The name of the saved search.")
712 filters: list[schemas.core.SavedSearchFilter] = Field(
713 default_factory=lambda: [], description="The filter set for the saved search."
714 )
717class ConcurrencyLimitCreate(ActionBaseModel):
718 """Data used by the Prefect REST API to create a concurrency limit."""
720 tag: str = Field(
721 default=..., description="A tag the concurrency limit is applied to."
722 )
723 concurrency_limit: int = Field(default=..., description="The concurrency limit.")
726class ConcurrencyLimitV2Create(ActionBaseModel):
727 """Data used by the Prefect REST API to create a v2 concurrency limit."""
729 active: bool = Field(
730 default=True, description="Whether the concurrency limit is active."
731 )
732 name: Name = Field(default=..., description="The name of the concurrency limit.")
733 limit: NonNegativeInteger = Field(default=..., description="The concurrency limit.")
734 active_slots: NonNegativeInteger = Field(
735 default=0, description="The number of active slots."
736 )
737 denied_slots: NonNegativeInteger = Field(
738 default=0, description="The number of denied slots."
739 )
740 slot_decay_per_second: NonNegativeFloat = Field(
741 default=0,
742 description="The decay rate for active slots when used as a rate limit.",
743 )
746class ConcurrencyLimitV2Update(ActionBaseModel):
747 """Data used by the Prefect REST API to update a v2 concurrency limit."""
749 active: Optional[bool] = Field(None)
750 name: Optional[Name] = Field(None)
751 limit: Optional[NonNegativeInteger] = Field(None)
752 active_slots: Optional[NonNegativeInteger] = Field(None)
753 denied_slots: Optional[NonNegativeInteger] = Field(None)
754 slot_decay_per_second: Optional[NonNegativeFloat] = Field(None)
757class BlockTypeCreate(ActionBaseModel):
758 """Data used by the Prefect REST API to create a block type."""
760 name: Name = Field(default=..., description="A block type's name")
761 slug: BlockTypeSlug = Field(default=..., description="A block type's slug")
762 logo_url: Optional[str] = Field( # TODO: HttpUrl
763 default=None, description="Web URL for the block type's logo"
764 )
765 documentation_url: Optional[str] = Field( # TODO: HttpUrl
766 default=None, description="Web URL for the block type's documentation"
767 )
768 description: Optional[str] = Field(
769 default=None,
770 description="A short blurb about the corresponding block's intended use",
771 )
772 code_example: Optional[str] = Field(
773 default=None,
774 description="A code snippet demonstrating use of the corresponding block",
775 )
778class BlockTypeUpdate(ActionBaseModel):
779 """Data used by the Prefect REST API to update a block type."""
781 logo_url: Optional[str] = Field(None) # TODO: HttpUrl
782 documentation_url: Optional[str] = Field(None) # TODO: HttpUrl
783 description: Optional[str] = Field(None)
784 code_example: Optional[str] = Field(None)
786 @classmethod
787 def updatable_fields(cls) -> set[str]:
788 return get_class_fields_only(cls)
791class BlockSchemaCreate(ActionBaseModel):
792 """Data used by the Prefect REST API to create a block schema."""
794 fields: dict[str, Any] = Field(
795 default_factory=dict, description="The block schema's field schema"
796 )
797 block_type_id: UUID = Field(default=..., description="A block type ID")
799 capabilities: List[str] = Field(
800 default_factory=list,
801 description="A list of Block capabilities",
802 )
803 version: str = Field(
804 default=schemas.core.DEFAULT_BLOCK_SCHEMA_VERSION,
805 description="Human readable identifier for the block schema",
806 )
809class BlockDocumentCreate(ActionBaseModel):
810 """Data used by the Prefect REST API to create a block document."""
812 name: Optional[BlockDocumentName] = Field(
813 default=None,
814 description=(
815 "The block document's name. Not required for anonymous block documents."
816 ),
817 )
818 data: Dict[str, Any] = Field(
819 default_factory=dict, description="The block document's data"
820 )
821 block_schema_id: UUID = Field(default=..., description="A block schema ID")
823 block_type_id: UUID = Field(default=..., description="A block type ID")
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)
840class BlockDocumentUpdate(ActionBaseModel):
841 """Data used by the Prefect REST API to update a block document."""
843 block_schema_id: Optional[UUID] = Field(
844 default=None, description="A block schema ID"
845 )
846 data: Dict[str, Any] = Field(
847 default_factory=dict, description="The block document's data"
848 )
849 merge_existing_data: bool = True
852class BlockDocumentReferenceCreate(ActionBaseModel):
853 """Data used to create block document reference."""
855 id: UUID = Field(
856 default_factory=uuid4, description="The block document reference ID"
857 )
858 parent_block_document_id: UUID = Field(
859 default=..., description="ID of the parent block document"
860 )
861 reference_block_document_id: UUID = Field(
862 default=..., description="ID of the nested block document"
863 )
864 name: str = Field(
865 default=..., description="The name that the reference is nested under"
866 )
868 @model_validator(mode="before")
869 def validate_parent_and_ref_are_different(cls, values):
870 return validate_parent_and_ref_diff(values)
873class LogCreate(ActionBaseModel):
874 """Data used by the Prefect REST API to create a log."""
876 name: str = Field(default=..., description="The logger name.")
877 level: int = Field(default=..., description="The log level.")
878 message: str = Field(default=..., description="The log message.")
879 timestamp: DateTime = Field(default=..., description="The log timestamp.")
880 flow_run_id: Optional[UUID] = Field(None)
881 task_run_id: Optional[UUID] = Field(None)
884def validate_base_job_template(v: dict[str, Any]) -> dict[str, Any]:
885 if v == dict():
886 return v
888 job_config = v.get("job_configuration")
889 variables_schema = v.get("variables")
890 if not (job_config and variables_schema): 890 ↛ 895line 890 didn't jump to line 895 because the condition on line 890 was always true
891 raise ValueError(
892 "The `base_job_template` must contain both a `job_configuration` key"
893 " and a `variables` key."
894 )
895 template_variables: set[str] = set()
896 for template in job_config.values():
897 # find any variables inside of double curly braces, minus any whitespace
898 # e.g. "{{ var1 }}.{{var2}}" -> ["var1", "var2"]
899 # convert to json string to handle nested objects and lists
900 found_variables = find_placeholders(json.dumps(template))
901 template_variables.update({placeholder.name for placeholder in found_variables})
903 provided_variables = set(variables_schema.get("properties", {}).keys())
904 if not template_variables.issubset(provided_variables):
905 missing_variables = template_variables - provided_variables
906 raise ValueError(
907 "The variables specified in the job configuration template must be "
908 "present as properties in the variables schema. "
909 "Your job configuration uses the following undeclared "
910 f"variable(s): {' ,'.join(missing_variables)}."
911 )
912 return v
915class WorkPoolCreate(ActionBaseModel):
916 """Data used by the Prefect REST API to create a work pool."""
918 name: NonEmptyishName = Field(..., description="The name of the work pool.")
919 description: Optional[str] = Field(
920 default=None, description="The work pool description."
921 )
922 type: str = Field(description="The work pool type.", default="prefect-agent")
923 base_job_template: Dict[str, Any] = Field(
924 default_factory=dict, description="The work pool's base job template."
925 )
926 is_paused: bool = Field(
927 default=False,
928 description="Pausing the work pool stops the delivery of all work.",
929 )
930 concurrency_limit: Optional[NonNegativeInteger] = Field(
931 default=None, description="A concurrency limit for the work pool."
932 )
934 storage_configuration: schemas.core.WorkPoolStorageConfiguration = Field(
935 default_factory=schemas.core.WorkPoolStorageConfiguration,
936 description="The storage configuration for the work pool.",
937 )
939 _validate_base_job_template = field_validator("base_job_template")(
940 validate_base_job_template
941 )
944class WorkPoolUpdate(ActionBaseModel):
945 """Data used by the Prefect REST API to update a work pool."""
947 description: Optional[str] = Field(default=None)
948 is_paused: Optional[bool] = Field(default=None)
949 base_job_template: Optional[Dict[str, Any]] = Field(default=None)
950 concurrency_limit: Optional[NonNegativeInteger] = Field(default=None)
951 storage_configuration: Optional[schemas.core.WorkPoolStorageConfiguration] = Field(
952 default=None,
953 description="The storage configuration for the work pool.",
954 )
955 _validate_base_job_template = field_validator("base_job_template")(
956 validate_base_job_template
957 )
960class WorkQueueCreate(ActionBaseModel):
961 """Data used by the Prefect REST API to create a work queue."""
963 name: Name = Field(default=..., description="The name of the work queue.")
964 description: Optional[str] = Field(
965 default="", description="An optional description for the work queue."
966 )
967 is_paused: bool = Field(
968 default=False, description="Whether or not the work queue is paused."
969 )
970 concurrency_limit: Optional[NonNegativeInteger] = Field(
971 None, description="The work queue's concurrency limit."
972 )
973 priority: Optional[PositiveInteger] = Field(
974 None,
975 description=(
976 "The queue's priority. Lower values are higher priority (1 is the highest)."
977 ),
978 )
980 # DEPRECATED
982 filter: Optional[schemas.core.QueueFilter] = Field(
983 None,
984 description="DEPRECATED: Filter criteria for the work queue.",
985 deprecated=True,
986 )
989class WorkQueueUpdate(ActionBaseModel):
990 """Data used by the Prefect REST API to update a work queue."""
992 name: Optional[str] = Field(None)
993 description: Optional[str] = Field(None)
994 is_paused: bool = Field(
995 default=False, description="Whether or not the work queue is paused."
996 )
997 concurrency_limit: Optional[NonNegativeInteger] = Field(None)
998 priority: Optional[PositiveInteger] = Field(None)
999 last_polled: Optional[DateTime] = Field(None)
1001 # DEPRECATED
1003 filter: Optional[schemas.core.QueueFilter] = Field(
1004 None,
1005 description="DEPRECATED: Filter criteria for the work queue.",
1006 deprecated=True,
1007 )
1010class ArtifactCreate(ActionBaseModel):
1011 """Data used by the Prefect REST API to create an artifact."""
1013 key: Optional[ArtifactKey] = Field(
1014 default=None, description="An optional unique reference key for this artifact."
1015 )
1016 type: Optional[str] = Field(
1017 default=None,
1018 description=(
1019 "An identifier that describes the shape of the data field. e.g. 'result',"
1020 " 'table', 'markdown'"
1021 ),
1022 )
1023 description: Optional[str] = Field(
1024 default=None, description="A markdown-enabled description of the artifact."
1025 )
1026 data: Optional[Union[Dict[str, Any], Any]] = Field(
1027 default=None,
1028 description=(
1029 "Data associated with the artifact, e.g. a result.; structure depends on"
1030 " the artifact type."
1031 ),
1032 )
1033 metadata_: Optional[
1034 Annotated[dict[str, str], AfterValidator(validate_max_metadata_length)]
1035 ] = Field(
1036 default=None,
1037 description=(
1038 "User-defined artifact metadata. Content must be string key and value"
1039 " pairs."
1040 ),
1041 )
1042 flow_run_id: Optional[UUID] = Field(
1043 default=None, description="The flow run associated with the artifact."
1044 )
1045 task_run_id: Optional[UUID] = Field(
1046 default=None, description="The task run associated with the artifact."
1047 )
1049 @classmethod
1050 def from_result(cls, data: Any | dict[str, Any]) -> "ArtifactCreate":
1051 artifact_info: dict[str, Any] = dict()
1052 if isinstance(data, dict):
1053 artifact_key = data.pop("artifact_key", None)
1054 if artifact_key:
1055 artifact_info["key"] = artifact_key
1057 artifact_type = data.pop("artifact_type", None)
1058 if artifact_type:
1059 artifact_info["type"] = artifact_type
1061 description = data.pop("artifact_description", None)
1062 if description:
1063 artifact_info["description"] = description
1065 return cls(data=data, **artifact_info)
1068class ArtifactUpdate(ActionBaseModel):
1069 """Data used by the Prefect REST API to update an artifact."""
1071 data: Optional[Union[Dict[str, Any], Any]] = Field(None)
1072 description: Optional[str] = Field(None)
1073 metadata_: Optional[
1074 Annotated[dict[str, str], AfterValidator(validate_max_metadata_length)]
1075 ] = Field(None)
1078class VariableCreate(ActionBaseModel):
1079 """Data used by the Prefect REST API to create a Variable."""
1081 name: VariableName = Field(default=...)
1082 value: StrictVariableValue = Field(
1083 default=...,
1084 description="The value of the variable",
1085 examples=["my-value"],
1086 )
1087 tags: list[str] = Field(
1088 default_factory=list,
1089 description="A list of variable tags",
1090 examples=[["tag-1", "tag-2"]],
1091 )
1094class VariableUpdate(ActionBaseModel):
1095 """Data used by the Prefect REST API to update a Variable."""
1097 name: Optional[VariableName] = Field(default=None)
1098 value: StrictVariableValue = Field(
1099 default=None,
1100 description="The value of the variable",
1101 examples=["my-value"],
1102 )
1103 tags: Optional[list[str]] = Field(
1104 default=None,
1105 description="A list of variable tags",
1106 examples=[["tag-1", "tag-2"]],
1107 )