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

1""" 

2Reduced schemas for accepting API actions. 

3""" 

4 

5from __future__ import annotations 

6 

7import json 

8from copy import deepcopy 

9from typing import Annotated, Any, ClassVar, Dict, List, Optional, Union 

10from uuid import UUID, uuid4 

11 

12from pydantic import ( 

13 AfterValidator, 

14 ConfigDict, 

15 Field, 

16 field_validator, 

17 model_validator, 

18) 

19 

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 

57 

58SizedParameters = Annotated[ 

59 Dict[str, Any], AfterValidator(validate_parameter_size_field) 

60] 

61 

62 

63class ActionBaseModel(PrefectBaseModel): 

64 model_config: ClassVar[ConfigDict] = ConfigDict(extra="forbid") 

65 

66 

67class FlowCreate(ActionBaseModel): 

68 """Data used by the Prefect REST API to create a flow.""" 

69 

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 ) 

83 

84 

85class FlowUpdate(ActionBaseModel): 

86 """Data used by the Prefect REST API to update a flow.""" 

87 

88 tags: List[str] = Field( 

89 default_factory=list, 

90 description="A list of flow tags", 

91 examples=[["tag-1", "tag-2"]], 

92 ) 

93 

94 

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] 

108 

109 

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 ) 

132 

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 ) 

141 

142 

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 ) 

150 

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 ) 

166 

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 ) 

175 

176 

177class DeploymentCreate(ActionBaseModel): 

178 """Data used by the Prefect REST API to create a deployment.""" 

179 

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) 

232 

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 ) 

250 

251 version_info: Optional[schemas.core.VersionInfo] = Field( 

252 default=None, description="A description of this version of the deployment." 

253 ) 

254 

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. 

259 

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 

267 

268 variables_schema = deepcopy(base_job_template.get("variables")) 

269 

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 ) 

278 

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) 

283 

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 

297 

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 

308 

309 

310class DeploymentUpdate(ActionBaseModel): 

311 """Data used by the Prefect REST API to update a deployment.""" 

312 

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) 

317 

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) 

372 

373 version_info: Optional[schemas.core.VersionInfo] = Field( 

374 default=None, description="A description of this version of the deployment." 

375 ) 

376 

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. 

381 

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 

389 

390 variables_schema = deepcopy(base_job_template.get("variables")) 

391 

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 

403 

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 

414 

415 

416class FlowRunUpdate(ActionBaseModel): 

417 """Data used by the Prefect REST API to update a flow run.""" 

418 

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) 

428 

429 @field_validator("name", mode="before") 

430 @classmethod 

431 def set_name(cls, name: str) -> str: 

432 return get_or_create_run_name(name) 

433 

434 

435class StateCreate(ActionBaseModel): 

436 """Data used by the Prefect REST API to create a new state.""" 

437 

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 ) 

454 

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 

464 

465 @model_validator(mode="after") 

466 def default_scheduled_start_time(self): 

467 from prefect.server.schemas.states import StateType 

468 

469 if self.type == StateType.SCHEDULED: 

470 if not self.state_details.scheduled_time: 

471 self.state_details.scheduled_time = now("UTC") 

472 

473 return self 

474 

475 

476class TaskRunCreate(ActionBaseModel): 

477 """Data used by the Prefect REST API to create a task run""" 

478 

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 ) 

487 

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 ) 

545 

546 @field_validator("name", mode="before") 

547 @classmethod 

548 def set_name(cls, name: str) -> str: 

549 return get_or_create_run_name(name) 

550 

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) 

555 

556 

557class TaskRunUpdate(ActionBaseModel): 

558 """Data used by the Prefect REST API to update a task run""" 

559 

560 name: str = Field( 

561 default_factory=lambda: generate_slug(2), examples=["my-task-run"] 

562 ) 

563 

564 @field_validator("name", mode="before") 

565 @classmethod 

566 def set_name(cls, name: str) -> str: 

567 return get_or_create_run_name(name) 

568 

569 

570class FlowRunCreate(ActionBaseModel): 

571 """Data used by the Prefect REST API to create a flow run.""" 

572 

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 ) 

577 

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 ) 

631 

632 # DEPRECATED 

633 

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 ) 

642 

643 @field_validator("name", mode="before") 

644 @classmethod 

645 def set_name(cls, name: str) -> str: 

646 return get_or_create_run_name(name) 

647 

648 

649class DeploymentFlowRunCreate(ActionBaseModel): 

650 """Data used by the Prefect REST API to create a flow run from a deployment.""" 

651 

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 ) 

656 

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 ) 

701 

702 @field_validator("name", mode="before") 

703 @classmethod 

704 def set_name(cls, name: str) -> str: 

705 return get_or_create_run_name(name) 

706 

707 

708class SavedSearchCreate(ActionBaseModel): 

709 """Data used by the Prefect REST API to create a saved search.""" 

710 

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 ) 

715 

716 

717class ConcurrencyLimitCreate(ActionBaseModel): 

718 """Data used by the Prefect REST API to create a concurrency limit.""" 

719 

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

724 

725 

726class ConcurrencyLimitV2Create(ActionBaseModel): 

727 """Data used by the Prefect REST API to create a v2 concurrency limit.""" 

728 

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 ) 

744 

745 

746class ConcurrencyLimitV2Update(ActionBaseModel): 

747 """Data used by the Prefect REST API to update a v2 concurrency limit.""" 

748 

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) 

755 

756 

757class BlockTypeCreate(ActionBaseModel): 

758 """Data used by the Prefect REST API to create a block type.""" 

759 

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 ) 

776 

777 

778class BlockTypeUpdate(ActionBaseModel): 

779 """Data used by the Prefect REST API to update a block type.""" 

780 

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) 

785 

786 @classmethod 

787 def updatable_fields(cls) -> set[str]: 

788 return get_class_fields_only(cls) 

789 

790 

791class BlockSchemaCreate(ActionBaseModel): 

792 """Data used by the Prefect REST API to create a block schema.""" 

793 

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

798 

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 ) 

807 

808 

809class BlockDocumentCreate(ActionBaseModel): 

810 """Data used by the Prefect REST API to create a block document.""" 

811 

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

822 

823 block_type_id: UUID = Field(default=..., description="A block type ID") 

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 ) 

832 

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) 

838 

839 

840class BlockDocumentUpdate(ActionBaseModel): 

841 """Data used by the Prefect REST API to update a block document.""" 

842 

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 

850 

851 

852class BlockDocumentReferenceCreate(ActionBaseModel): 

853 """Data used to create block document reference.""" 

854 

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 ) 

867 

868 @model_validator(mode="before") 

869 def validate_parent_and_ref_are_different(cls, values): 

870 return validate_parent_and_ref_diff(values) 

871 

872 

873class LogCreate(ActionBaseModel): 

874 """Data used by the Prefect REST API to create a log.""" 

875 

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) 

882 

883 

884def validate_base_job_template(v: dict[str, Any]) -> dict[str, Any]: 

885 if v == dict(): 

886 return v 

887 

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

902 

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 

913 

914 

915class WorkPoolCreate(ActionBaseModel): 

916 """Data used by the Prefect REST API to create a work pool.""" 

917 

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 ) 

933 

934 storage_configuration: schemas.core.WorkPoolStorageConfiguration = Field( 

935 default_factory=schemas.core.WorkPoolStorageConfiguration, 

936 description="The storage configuration for the work pool.", 

937 ) 

938 

939 _validate_base_job_template = field_validator("base_job_template")( 

940 validate_base_job_template 

941 ) 

942 

943 

944class WorkPoolUpdate(ActionBaseModel): 

945 """Data used by the Prefect REST API to update a work pool.""" 

946 

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 ) 

958 

959 

960class WorkQueueCreate(ActionBaseModel): 

961 """Data used by the Prefect REST API to create a work queue.""" 

962 

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 ) 

979 

980 # DEPRECATED 

981 

982 filter: Optional[schemas.core.QueueFilter] = Field( 

983 None, 

984 description="DEPRECATED: Filter criteria for the work queue.", 

985 deprecated=True, 

986 ) 

987 

988 

989class WorkQueueUpdate(ActionBaseModel): 

990 """Data used by the Prefect REST API to update a work queue.""" 

991 

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) 

1000 

1001 # DEPRECATED 

1002 

1003 filter: Optional[schemas.core.QueueFilter] = Field( 

1004 None, 

1005 description="DEPRECATED: Filter criteria for the work queue.", 

1006 deprecated=True, 

1007 ) 

1008 

1009 

1010class ArtifactCreate(ActionBaseModel): 

1011 """Data used by the Prefect REST API to create an artifact.""" 

1012 

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 ) 

1048 

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 

1056 

1057 artifact_type = data.pop("artifact_type", None) 

1058 if artifact_type: 

1059 artifact_info["type"] = artifact_type 

1060 

1061 description = data.pop("artifact_description", None) 

1062 if description: 

1063 artifact_info["description"] = description 

1064 

1065 return cls(data=data, **artifact_info) 

1066 

1067 

1068class ArtifactUpdate(ActionBaseModel): 

1069 """Data used by the Prefect REST API to update an artifact.""" 

1070 

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) 

1076 

1077 

1078class VariableCreate(ActionBaseModel): 

1079 """Data used by the Prefect REST API to create a Variable.""" 

1080 

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 ) 

1092 

1093 

1094class VariableUpdate(ActionBaseModel): 

1095 """Data used by the Prefect REST API to update a Variable.""" 

1096 

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 )