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

1""" 

2Full schemas of Prefect REST API objects. 

3""" 

4 

5from __future__ import annotations 

6 

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 

20 

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 

34 

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 

74 

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 

77 

78DEFAULT_BLOCK_SCHEMA_VERSION = "non-versioned" 

79 

80KeyValueLabels = dict[str, Union[StrictBool, StrictInt, StrictFloat, str]] 

81 

82 

83class Flow(ORMBaseModel): 

84 """An ORM representation of flow data.""" 

85 

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 ) 

99 

100 

101class FlowRunPolicy(PrefectBaseModel): 

102 """Defines of how a flow run should retry.""" 

103 

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 ) 

133 

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) 

137 

138 

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 ) 

149 

150 

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 ) 

161 

162 

163class ConcurrencyLimitStrategy(AutoEnum): 

164 """ 

165 Enumeration of concurrency collision strategies. 

166 """ 

167 

168 ENQUEUE = AutoEnum.auto() 

169 CANCEL_NEW = AutoEnum.auto() 

170 

171 

172class ConcurrencyOptions(BaseModel): 

173 """ 

174 Class for storing the concurrency config in database. 

175 """ 

176 

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 ) 

184 

185 

186class FlowRun(TimeSeriesBaseModel, ORMBaseModel): 

187 """An ORM representation of flow run data.""" 

188 

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 ) 

254 

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 ) 

312 

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 

320 

321 job_variables: Optional[Dict[str, Any]] = Field( 

322 default=None, 

323 description="Variables used as overrides in the base job template", 

324 ) 

325 

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

327 @classmethod 

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

329 return get_or_create_run_name(name) 

330 

331 def __eq__(self, other: Any) -> bool: 

332 """ 

333 Check for "equality" to another flow run schema 

334 

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) 

344 

345 

346class TaskRunPolicy(PrefectBaseModel): 

347 """Defines of how a task run should retry.""" 

348 

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 ) 

373 

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) 

377 

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) 

384 

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) 

389 

390 

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

396 

397 model_config: ClassVar[ConfigDict] = ConfigDict(frozen=True) 

398 

399 input_type: str 

400 

401 

402class TaskRunResult(RunInput): 

403 """Represents a task run result input to another task run.""" 

404 

405 input_type: Literal["task_run"] = "task_run" 

406 id: UUID 

407 

408 

409class FlowRunResult(RunInput): 

410 input_type: Literal["flow_run"] = "flow_run" 

411 id: UUID 

412 

413 

414class Parameter(RunInput): 

415 """Represents a parameter input to a task run.""" 

416 

417 input_type: Literal["parameter"] = "parameter" 

418 name: str 

419 

420 

421class Constant(RunInput): 

422 """Represents constant input value to a task run.""" 

423 

424 input_type: Literal["constant"] = "constant" 

425 type: str 

426 

427 

428class TaskRun(TimeSeriesBaseModel, ORMBaseModel): 

429 """An ORM representation of task run data.""" 

430 

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 ) 

505 

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 ) 

533 

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 ) 

540 

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

542 @classmethod 

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

544 return get_or_create_run_name(name) 

545 

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) 

550 

551 

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 ) 

574 

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 ) 

581 

582 

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

586 

587 

588class Deployment(ORMBaseModel): 

589 """An ORM representation of deployment data.""" 

590 

591 model_config: ClassVar[ConfigDict] = ConfigDict(populate_by_name=True) 

592 

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 ) 

698 

699 

700class ConcurrencyLimit(ORMBaseModel): 

701 """An ORM representation of a concurrency limit.""" 

702 

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 ) 

711 

712 

713class ConcurrencyLimitV2(ORMBaseModel): 

714 """An ORM representation of a v2 concurrency limit.""" 

715 

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 ) 

730 

731 

732class BlockType(ORMBaseModel): 

733 """An ORM representation of a block type""" 

734 

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 ) 

754 

755 

756class BlockSchema(ORMBaseModel): 

757 """An ORM representation of a block schema.""" 

758 

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 ) 

777 

778 

779class BlockSchemaReference(ORMBaseModel): 

780 """An ORM representation of a block schema reference.""" 

781 

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 ) 

797 

798 

799class BlockDocument(ORMBaseModel): 

800 """An ORM representation of a block document.""" 

801 

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 ) 

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

885 

886 

887class BlockDocumentReference(ORMBaseModel): 

888 """An ORM representation of a block document reference.""" 

889 

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 ) 

905 

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) 

911 

912 

913class Configuration(ORMBaseModel): 

914 """An ORM representation of account info.""" 

915 

916 key: str = Field(default=..., description="Account info key") 

917 value: Dict[str, Any] = Field(default=..., description="Account info") 

918 

919 

920class ServerDefaultResultStorage(PrefectBaseModel): 

921 """Server-side default result storage configuration.""" 

922 

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 ) 

927 

928 

929class ServerDefaultResultStorageUpdate(PrefectBaseModel): 

930 """Request payload for setting the server default result storage block.""" 

931 

932 default_result_storage_block_id: UUID = Field( 

933 default=..., 

934 description="The block document ID of the server default result storage block.", 

935 ) 

936 

937 

938class SavedSearchFilter(PrefectBaseModel): 

939 """A filter for a saved search model. Intended for use by the Prefect UI.""" 

940 

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 ) 

953 

954 

955class SavedSearch(ORMBaseModel): 

956 """An ORM representation of saved search data. Represents a set of filter criteria.""" 

957 

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 ) 

963 

964 

965class Log(TimeSeriesBaseModel, ORMBaseModel): 

966 """An ORM representation of log data.""" 

967 

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 ) 

978 

979 

980class QueueFilter(PrefectBaseModel): 

981 """Filter criteria definition for a work queue.""" 

982 

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 ) 

991 

992 

993class WorkQueue(ORMBaseModel): 

994 """An ORM representation of a work queue""" 

995 

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 ) 

1024 

1025 

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 ) 

1041 

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. 

1047 

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. 

1051 

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 

1061 

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 

1069 

1070 return healthy 

1071 

1072 

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 ) 

1087 

1088 

1089class Agent(ORMBaseModel): 

1090 """An ORM representation of an agent""" 

1091 

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 ) 

1105 

1106 

1107class WorkPoolStorageConfiguration(PrefectBaseModel): 

1108 """A representation of a work pool's storage configuration""" 

1109 

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

1111 

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 ) 

1124 

1125 

1126class WorkPool(ORMBaseModel): 

1127 """An ORM representation of a work pool""" 

1128 

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 ) 

1149 

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 ) 

1155 

1156 storage_configuration: WorkPoolStorageConfiguration = Field( 

1157 default_factory=WorkPoolStorageConfiguration, 

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

1159 ) 

1160 

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) 

1164 

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 

1181 

1182 

1183class Worker(ORMBaseModel): 

1184 """An ORM representation of a worker""" 

1185 

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 ) 

1199 

1200 

1201Flow.model_rebuild() 

1202FlowRun.model_rebuild() 

1203 

1204 

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 ) 

1240 

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 

1248 

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 

1252 

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 

1256 

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

1258 

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) 

1263 

1264 

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 ) 

1300 

1301 

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 ) 

1319 

1320 

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

1328 

1329 

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 )