Coverage for /usr/local/lib/python3.12/site-packages/prefect/server/schemas/responses.py: 87%

280 statements  

« prev     ^ index     » next       coverage.py v7.15.2, created at 2026-10-07 02:04 +0000

1""" 

2Schemas for special responses from the Prefect REST API. 

3""" 

4 

5import datetime 

6from datetime import timedelta 

7from typing import Any, ClassVar, Dict, List, Optional, Type, Union 

8from uuid import UUID 

9 

10from pydantic import BaseModel, ConfigDict, Field, model_validator 

11from typing_extensions import Literal, Self 

12 

13import prefect.server.schemas as schemas 

14from prefect.server.schemas.core import ( 

15 CreatedBy, 

16 FlowRunPolicy, 

17 UpdatedBy, 

18 WorkQueueStatusDetail, 

19) 

20from prefect.server.utilities.schemas.bases import ORMBaseModel, PrefectBaseModel 

21from prefect.types import DateTime, KeyValueLabelsField 

22from prefect.types._datetime import create_datetime_instance 

23from prefect.utilities.collections import AutoEnum 

24from prefect.utilities.names import generate_slug 

25 

26 

27class SetStateStatus(AutoEnum): 

28 """Enumerates return statuses for setting run states.""" 

29 

30 ACCEPT = AutoEnum.auto() 

31 REJECT = AutoEnum.auto() 

32 ABORT = AutoEnum.auto() 

33 WAIT = AutoEnum.auto() 

34 

35 

36class StateAcceptDetails(PrefectBaseModel): 

37 """Details associated with an ACCEPT state transition.""" 

38 

39 type: Literal["accept_details"] = Field( 

40 default="accept_details", 

41 description=( 

42 "The type of state transition detail. Used to ensure pydantic does not" 

43 " coerce into a different type." 

44 ), 

45 ) 

46 

47 

48class StateRejectDetails(PrefectBaseModel): 

49 """Details associated with a REJECT state transition.""" 

50 

51 type: Literal["reject_details"] = Field( 

52 default="reject_details", 

53 description=( 

54 "The type of state transition detail. Used to ensure pydantic does not" 

55 " coerce into a different type." 

56 ), 

57 ) 

58 reason: Optional[str] = Field( 

59 default=None, description="The reason why the state transition was rejected." 

60 ) 

61 

62 

63class StateAbortDetails(PrefectBaseModel): 

64 """Details associated with an ABORT state transition.""" 

65 

66 type: Literal["abort_details"] = Field( 

67 default="abort_details", 

68 description=( 

69 "The type of state transition detail. Used to ensure pydantic does not" 

70 " coerce into a different type." 

71 ), 

72 ) 

73 reason: Optional[str] = Field( 

74 default=None, description="The reason why the state transition was aborted." 

75 ) 

76 

77 

78class StateWaitDetails(PrefectBaseModel): 

79 """Details associated with a WAIT state transition.""" 

80 

81 type: Literal["wait_details"] = Field( 

82 default="wait_details", 

83 description=( 

84 "The type of state transition detail. Used to ensure pydantic does not" 

85 " coerce into a different type." 

86 ), 

87 ) 

88 delay_seconds: int = Field( 

89 default=..., 

90 description=( 

91 "The length of time in seconds the client should wait before transitioning" 

92 " states." 

93 ), 

94 ) 

95 reason: Optional[str] = Field( 

96 default=None, description="The reason why the state transition should wait." 

97 ) 

98 

99 

100class HistoryResponseState(PrefectBaseModel): 

101 """Represents a single state's history over an interval.""" 

102 

103 state_type: schemas.states.StateType = Field( 

104 default=..., description="The state type." 

105 ) 

106 state_name: str = Field(default=..., description="The state name.") 

107 count_runs: int = Field( 

108 default=..., 

109 description="The number of runs in the specified state during the interval.", 

110 ) 

111 sum_estimated_run_time: datetime.timedelta = Field( 

112 default=..., 

113 description="The total estimated run time of all runs during the interval.", 

114 ) 

115 sum_estimated_lateness: datetime.timedelta = Field( 

116 default=..., 

117 description=( 

118 "The sum of differences between actual and expected start time during the" 

119 " interval." 

120 ), 

121 ) 

122 

123 

124class HistoryResponse(PrefectBaseModel): 

125 """Represents a history of aggregation states over an interval""" 

126 

127 interval_start: DateTime = Field( 

128 default=..., description="The start date of the interval." 

129 ) 

130 interval_end: DateTime = Field( 

131 default=..., description="The end date of the interval." 

132 ) 

133 states: List[HistoryResponseState] = Field( 

134 default=..., description="A list of state histories during the interval." 

135 ) 

136 

137 @model_validator(mode="before") 

138 @classmethod 

139 def validate_timestamps( 

140 cls, values: dict 

141 ) -> dict: # TODO: remove this, handle with ORM 

142 d = {"interval_start": None, "interval_end": None} 

143 for field in d.keys(): 

144 val = values.get(field) 

145 if isinstance(val, datetime.datetime): 

146 d[field] = create_datetime_instance(values[field]) 

147 else: 

148 d[field] = val 

149 

150 return {**values, **d} 

151 

152 

153StateResponseDetails = Union[ 

154 StateAcceptDetails, StateWaitDetails, StateRejectDetails, StateAbortDetails 

155] 

156 

157 

158class OrchestrationResult(PrefectBaseModel): 

159 """ 

160 A container for the output of state orchestration. 

161 """ 

162 

163 state: Optional[schemas.states.State] 

164 status: SetStateStatus 

165 details: StateResponseDetails 

166 

167 

168class WorkerFlowRunResponse(PrefectBaseModel): 

169 model_config: ClassVar[ConfigDict] = ConfigDict(arbitrary_types_allowed=True) 

170 

171 work_pool_id: UUID 

172 work_queue_id: UUID 

173 flow_run: schemas.core.FlowRun 

174 

175 

176class FlowRunResponse(ORMBaseModel): 

177 name: str = Field( 

178 default_factory=lambda: generate_slug(2), 

179 description=( 

180 "The name of the flow run. Defaults to a random slug if not specified." 

181 ), 

182 examples=["my-flow-run"], 

183 ) 

184 flow_id: UUID = Field(default=..., description="The id of the flow being run.") 

185 state_id: Optional[UUID] = Field( 

186 default=None, description="The id of the flow run's current state." 

187 ) 

188 deployment_id: Optional[UUID] = Field( 

189 default=None, 

190 description=( 

191 "The id of the deployment associated with this flow run, if available." 

192 ), 

193 ) 

194 deployment_version: Optional[str] = Field( 

195 default=None, 

196 description="The version of the deployment associated with this flow run.", 

197 examples=["1.0"], 

198 ) 

199 work_queue_id: Optional[UUID] = Field( 

200 default=None, description="The id of the run's work pool queue." 

201 ) 

202 work_queue_name: Optional[str] = Field( 

203 default=None, description="The work queue that handled this flow run." 

204 ) 

205 flow_version: Optional[str] = Field( 

206 default=None, 

207 description="The version of the flow executed in this flow run.", 

208 examples=["1.0"], 

209 ) 

210 parameters: Dict[str, Any] = Field( 

211 default_factory=dict, description="Parameters for the flow run." 

212 ) 

213 idempotency_key: Optional[str] = Field( 

214 default=None, 

215 description=( 

216 "An optional idempotency key for the flow run. Used to ensure the same flow" 

217 " run is not created multiple times." 

218 ), 

219 ) 

220 context: Dict[str, Any] = Field( 

221 default_factory=dict, 

222 description="Additional context for the flow run.", 

223 examples=[{"my_var": "my_val"}], 

224 ) 

225 empirical_policy: FlowRunPolicy = Field( 

226 default_factory=FlowRunPolicy, 

227 ) 

228 tags: List[str] = Field( 

229 default_factory=list, 

230 description="A list of tags on the flow run", 

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

232 ) 

233 labels: KeyValueLabelsField 

234 parent_task_run_id: Optional[UUID] = Field( 

235 default=None, 

236 description=( 

237 "If the flow run is a subflow, the id of the 'dummy' task in the parent" 

238 " flow used to track subflow state." 

239 ), 

240 ) 

241 state_type: Optional[schemas.states.StateType] = Field( 

242 default=None, description="The type of the current flow run state." 

243 ) 

244 state_name: Optional[str] = Field( 

245 default=None, description="The name of the current flow run state." 

246 ) 

247 run_count: int = Field( 

248 default=0, description="The number of times the flow run was executed." 

249 ) 

250 expected_start_time: Optional[DateTime] = Field( 

251 default=None, 

252 description="The flow run's expected start time.", 

253 ) 

254 next_scheduled_start_time: Optional[DateTime] = Field( 

255 default=None, 

256 description="The next time the flow run is scheduled to start.", 

257 ) 

258 start_time: Optional[DateTime] = Field( 

259 default=None, description="The actual start time." 

260 ) 

261 end_time: Optional[DateTime] = Field( 

262 default=None, description="The actual end time." 

263 ) 

264 total_run_time: datetime.timedelta = Field( 

265 default=datetime.timedelta(0), 

266 description=( 

267 "Total run time. If the flow run was executed multiple times, the time of" 

268 " each run will be summed." 

269 ), 

270 ) 

271 estimated_run_time: datetime.timedelta = Field( 

272 default=datetime.timedelta(0), 

273 description="A real-time estimate of the total run time.", 

274 ) 

275 estimated_start_time_delta: datetime.timedelta = Field( 

276 default=datetime.timedelta(0), 

277 description="The difference between actual and expected start time.", 

278 ) 

279 auto_scheduled: bool = Field( 

280 default=False, 

281 description="Whether or not the flow run was automatically scheduled.", 

282 ) 

283 infrastructure_document_id: Optional[UUID] = Field( 

284 default=None, 

285 description="The block document defining infrastructure to use this flow run.", 

286 ) 

287 infrastructure_pid: Optional[str] = Field( 

288 default=None, 

289 description="The id of the flow run as returned by an infrastructure block.", 

290 ) 

291 created_by: Optional[CreatedBy] = Field( 

292 default=None, 

293 description="Optional information about the creator of this flow run.", 

294 ) 

295 work_pool_id: Optional[UUID] = Field( 

296 default=None, 

297 description="The id of the flow run's work pool.", 

298 ) 

299 work_pool_name: Optional[str] = Field( 

300 default=None, 

301 description="The name of the flow run's work pool.", 

302 examples=["my-work-pool"], 

303 ) 

304 state: Optional[schemas.states.State] = Field( 

305 default=None, description="The current state of the flow run." 

306 ) 

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

308 default=None, 

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

310 ) 

311 

312 @classmethod 

313 def model_validate( 

314 cls: Type[Self], 

315 obj: Any, 

316 *, 

317 strict: Optional[bool] = None, 

318 from_attributes: Optional[bool] = None, 

319 context: Optional[dict[str, Any]] = None, 

320 ) -> Self: 

321 response = super().model_validate(obj) 

322 

323 if from_attributes: 323 ↛ 331line 323 didn't jump to line 331 because the condition on line 323 was always true

324 if obj.work_queue: 

325 response.work_queue_id = obj.work_queue.id 

326 response.work_queue_name = obj.work_queue.name 

327 if obj.work_queue.work_pool: 327 ↛ 331line 327 didn't jump to line 331 because the condition on line 327 was always true

328 response.work_pool_id = obj.work_queue.work_pool.id 

329 response.work_pool_name = obj.work_queue.work_pool.name 

330 

331 return response 

332 

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

334 """ 

335 Check for "equality" to another flow run schema 

336 

337 Estimates times are rolling and will always change with repeated queries for 

338 a flow run so we ignore them during equality checks. 

339 """ 

340 if isinstance(other, FlowRunResponse): 

341 exclude_fields = {"estimated_run_time", "estimated_start_time_delta"} 

342 return self.model_dump(exclude=exclude_fields) == other.model_dump( 

343 exclude=exclude_fields 

344 ) 

345 return super().__eq__(other) 

346 

347 

348class TaskRunResponse(ORMBaseModel): 

349 name: str = Field( 

350 default_factory=lambda: generate_slug(2), 

351 description=( 

352 "The name of the task run. Defaults to a random slug if not specified." 

353 ), 

354 examples=["my-task-run"], 

355 ) 

356 flow_run_id: Optional[UUID] = Field( 

357 default=None, description="The id of the flow run this task run belongs to." 

358 ) 

359 task_key: str = Field( 

360 default=..., description="The key of the task this run represents." 

361 ) 

362 state_id: Optional[UUID] = Field( 

363 default=None, description="The id of the task run's current state." 

364 ) 

365 state: Optional[schemas.states.State] = Field( 

366 default=None, description="The current state of the task run." 

367 ) 

368 task_version: Optional[str] = Field( 

369 default=None, 

370 description="The version of the task executed in this task run.", 

371 examples=["1.0"], 

372 ) 

373 task_inputs: dict[ 

374 str, 

375 list[ 

376 Union[ 

377 schemas.core.TaskRunResult, 

378 schemas.core.FlowRunResult, 

379 schemas.core.Parameter, 

380 schemas.core.Constant, 

381 ] 

382 ], 

383 ] = Field(default_factory=dict, description="Inputs provided to the task run.") 

384 empirical_policy: schemas.core.TaskRunPolicy = Field( 

385 default_factory=schemas.core.TaskRunPolicy, 

386 description="The task run's empirical retry policy.", 

387 ) 

388 tags: list[str] = Field( 

389 default_factory=list, 

390 description="A list of tags for the task run.", 

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

392 ) 

393 start_time: Optional[DateTime] = Field( 

394 default=None, description="The actual start time." 

395 ) 

396 end_time: Optional[DateTime] = Field( 

397 default=None, description="The actual end time." 

398 ) 

399 total_run_time: datetime.timedelta = Field( 

400 default=datetime.timedelta(0), 

401 description=( 

402 "Total run time. If the task run was executed multiple times, the time of" 

403 " each run will be summed." 

404 ), 

405 ) 

406 estimated_run_time: datetime.timedelta = Field( 

407 default=datetime.timedelta(0), 

408 description="A real-time estimate of the total run time.", 

409 ) 

410 estimated_start_time_delta: datetime.timedelta = Field( 

411 default=datetime.timedelta(0), 

412 description="The difference between actual and expected start time.", 

413 ) 

414 

415 

416class DeploymentResponse(ORMBaseModel): 

417 name: str = Field(default=..., description="The name of the deployment.") 

418 version: Optional[str] = Field( 

419 default=None, description="An optional version for the deployment." 

420 ) 

421 description: Optional[str] = Field( 

422 default=None, description="A description for the deployment." 

423 ) 

424 flow_id: UUID = Field( 

425 default=..., description="The flow id associated with the deployment." 

426 ) 

427 paused: bool = Field( 

428 default=False, description="Whether or not the deployment is paused." 

429 ) 

430 schedules: List[schemas.core.DeploymentSchedule] = Field( 

431 default_factory=list, description="A list of schedules for the deployment." 

432 ) 

433 concurrency_limit: Optional[int] = Field( 

434 default=None, 

435 description="DEPRECATED: Prefer `global_concurrency_limit`. Will always be None for backwards compatibility. Will be removed after December 2024.", 

436 deprecated=True, 

437 ) 

438 global_concurrency_limit: Optional["GlobalConcurrencyLimitResponse"] = Field( 

439 default=None, 

440 description="The global concurrency limit object for enforcing the maximum number of flow runs that can be active at once.", 

441 ) 

442 concurrency_options: Optional[schemas.core.ConcurrencyOptions] = Field( 

443 default=None, 

444 description="The concurrency options for the deployment.", 

445 ) 

446 job_variables: Dict[str, Any] = Field( 

447 default_factory=dict, 

448 description="Overrides to apply to the base infrastructure block at runtime.", 

449 ) 

450 parameters: Dict[str, Any] = Field( 

451 default_factory=dict, 

452 description="Parameters for flow runs scheduled by the deployment.", 

453 ) 

454 tags: List[str] = Field( 

455 default_factory=list, 

456 description="A list of tags for the deployment", 

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

458 ) 

459 labels: KeyValueLabelsField 

460 work_queue_name: Optional[str] = Field( 

461 default=None, 

462 description=( 

463 "The work queue for the deployment. If no work queue is set, work will not" 

464 " be scheduled." 

465 ), 

466 ) 

467 work_queue_id: Optional[UUID] = Field( 

468 default=None, 

469 description="The id of the work pool queue to which this deployment is assigned.", 

470 ) 

471 last_polled: Optional[DateTime] = Field( 

472 default=None, 

473 description="The last time the deployment was polled for status updates.", 

474 ) 

475 parameter_openapi_schema: Optional[Dict[str, Any]] = Field( 

476 default=None, 

477 description="The parameter schema of the flow, including defaults.", 

478 json_schema_extra={"additionalProperties": True}, 

479 ) 

480 path: Optional[str] = Field( 

481 default=None, 

482 description=( 

483 "The path to the working directory for the workflow, relative to remote" 

484 " storage or an absolute path." 

485 ), 

486 ) 

487 pull_steps: Optional[list[dict[str, Any]]] = Field( 

488 default=None, description="Pull steps for cloning and running this deployment." 

489 ) 

490 entrypoint: Optional[str] = Field( 

491 default=None, 

492 description=( 

493 "The path to the entrypoint for the workflow, relative to the `path`." 

494 ), 

495 ) 

496 storage_document_id: Optional[UUID] = Field( 

497 default=None, 

498 description="The block document defining storage used for this flow.", 

499 ) 

500 infrastructure_document_id: Optional[UUID] = Field( 

501 default=None, 

502 description="The block document defining infrastructure to use for flow runs.", 

503 ) 

504 created_by: Optional[CreatedBy] = Field( 

505 default=None, 

506 description="Optional information about the creator of this deployment.", 

507 ) 

508 updated_by: Optional[UpdatedBy] = Field( 

509 default=None, 

510 description="Optional information about the updater of this deployment.", 

511 ) 

512 work_pool_name: Optional[str] = Field( 

513 default=None, description="The name of the deployment's work pool." 

514 ) 

515 status: Optional[schemas.statuses.DeploymentStatus] = Field( 

516 default=schemas.statuses.DeploymentStatus.NOT_READY, 

517 description="Whether the deployment is ready to run flows.", 

518 ) 

519 enforce_parameter_schema: bool = Field( 

520 default=True, 

521 description=( 

522 "Whether or not the deployment should enforce the parameter schema." 

523 ), 

524 ) 

525 

526 @classmethod 

527 def model_validate( 

528 cls: Type[Self], 

529 obj: Any, 

530 *, 

531 strict: Optional[bool] = None, 

532 from_attributes: Optional[bool] = None, 

533 context: Optional[dict[str, Any]] = None, 

534 ) -> Self: 

535 response = super().model_validate( 

536 obj, strict=strict, from_attributes=from_attributes, context=context 

537 ) 

538 

539 if from_attributes: 539 ↛ 546line 539 didn't jump to line 546 because the condition on line 539 was always true

540 if obj.work_queue: 540 ↛ 541line 540 didn't jump to line 541 because the condition on line 540 was never true

541 response.work_queue_id = obj.work_queue.id 

542 response.work_queue_name = obj.work_queue.name 

543 if obj.work_queue.work_pool: 

544 response.work_pool_name = obj.work_queue.work_pool.name 

545 

546 return response 

547 

548 

549class WorkPoolResponse(schemas.core.WorkPool): 

550 active_slots: Optional[int] = Field( 

551 default=None, 

552 description=( 

553 "The number of concurrency slots occupied by pending or running " 

554 "flow runs. None when concurrency_limit is not set." 

555 ), 

556 ) 

557 

558 

559class WorkQueueResponse(schemas.core.WorkQueue): 

560 work_pool_name: Optional[str] = Field( 

561 default=None, 

562 description="The name of the work pool the work pool resides within.", 

563 ) 

564 status: Optional[schemas.statuses.WorkQueueStatus] = Field( 

565 default=None, description="The queue status." 

566 ) 

567 active_slots: Optional[int] = Field( 

568 default=None, 

569 description=( 

570 "The number of concurrency slots currently in use. " 

571 "None when concurrency_limit is not set." 

572 ), 

573 ) 

574 

575 @classmethod 

576 def model_validate( 

577 cls: Type[Self], 

578 obj: Any, 

579 *, 

580 strict: Optional[bool] = None, 

581 from_attributes: Optional[bool] = None, 

582 context: Optional[dict[str, Any]] = None, 

583 ) -> Self: 

584 response = super().model_validate( 

585 obj, strict=strict, from_attributes=from_attributes, context=context 

586 ) 

587 

588 if from_attributes: 588 ↛ 592line 588 didn't jump to line 592 because the condition on line 588 was always true

589 if obj.work_pool: 589 ↛ 592line 589 didn't jump to line 592 because the condition on line 589 was always true

590 response.work_pool_name = obj.work_pool.name 

591 

592 return response 

593 

594 

595class WorkQueueWithStatus(WorkQueueResponse, WorkQueueStatusDetail): 

596 """Combines a work queue and its status details into a single object""" 

597 

598 

599DEFAULT_HEARTBEAT_INTERVAL_SECONDS = 30 

600INACTIVITY_HEARTBEAT_MULTIPLE = 3 

601 

602 

603class WorkerResponse(schemas.core.Worker): 

604 status: schemas.statuses.WorkerStatus = Field( 

605 schemas.statuses.WorkerStatus.OFFLINE, 

606 description="Current status of the worker.", 

607 ) 

608 

609 @classmethod 

610 def model_validate( 

611 cls: Type[Self], 

612 obj: Any, 

613 *, 

614 strict: Optional[bool] = None, 

615 from_attributes: Optional[bool] = None, 

616 context: Optional[dict[str, Any]] = None, 

617 ) -> Self: 

618 worker = super().model_validate( 

619 obj, strict=strict, from_attributes=from_attributes, context=context 

620 ) 

621 

622 if from_attributes: 

623 offline_horizon = datetime.datetime.now( 

624 tz=datetime.timezone.utc 

625 ) - datetime.timedelta( 

626 seconds=( 

627 worker.heartbeat_interval_seconds 

628 or DEFAULT_HEARTBEAT_INTERVAL_SECONDS 

629 ) 

630 * INACTIVITY_HEARTBEAT_MULTIPLE 

631 ) 

632 if worker.last_heartbeat_time > offline_horizon: 

633 worker.status = schemas.statuses.WorkerStatus.ONLINE 

634 else: 

635 worker.status = schemas.statuses.WorkerStatus.OFFLINE 

636 

637 return worker 

638 

639 

640class GlobalConcurrencyLimitResponse(ORMBaseModel): 

641 """ 

642 A response object for global concurrency limits. 

643 """ 

644 

645 active: bool = Field( 

646 default=True, description="Whether the global concurrency limit is active." 

647 ) 

648 name: str = Field( 

649 default=..., description="The name of the global concurrency limit." 

650 ) 

651 limit: int = Field(default=..., description="The concurrency limit.") 

652 active_slots: int = Field(default=..., description="The number of active slots.") 

653 slot_decay_per_second: float = Field( 

654 default=2.0, 

655 description="The decay rate for active slots when used as a rate limit.", 

656 ) 

657 

658 

659class FlowPaginationResponse(BaseModel): 

660 results: list[schemas.core.Flow] 

661 count: int 

662 limit: int 

663 pages: int 

664 page: int 

665 

666 

667class FlowRunPaginationResponse(BaseModel): 

668 results: list[FlowRunResponse] 

669 count: int 

670 limit: int 

671 pages: int 

672 page: int 

673 

674 

675class TaskRunPaginationResponse(BaseModel): 

676 results: list[TaskRunResponse] 

677 count: int 

678 limit: int 

679 pages: int 

680 page: int 

681 

682 

683class DeploymentPaginationResponse(BaseModel): 

684 results: list[DeploymentResponse] 

685 count: int 

686 limit: int 

687 pages: int 

688 page: int 

689 

690 

691class SchemaValuePropertyError(BaseModel): 

692 property: str 

693 errors: List["SchemaValueError"] 

694 

695 

696class SchemaValueIndexError(BaseModel): 

697 index: int 

698 errors: List["SchemaValueError"] 

699 

700 

701SchemaValueError = Union[str, SchemaValuePropertyError, SchemaValueIndexError] 

702 

703 

704class SchemaValuesValidationResponse(BaseModel): 

705 errors: List[SchemaValueError] 

706 valid: bool 

707 

708 

709# Bulk operation response schemas 

710 

711 

712class FlowRunBulkDeleteResponse(PrefectBaseModel): 

713 """Response from bulk flow run deletion.""" 

714 

715 deleted: List[UUID] = Field(default_factory=list) 

716 

717 

718class DeploymentBulkDeleteResponse(PrefectBaseModel): 

719 """Response from bulk deployment deletion.""" 

720 

721 deleted: List[UUID] = Field(default_factory=list) 

722 

723 

724class FlowBulkDeleteResponse(PrefectBaseModel): 

725 """Response from bulk flow deletion.""" 

726 

727 deleted: List[UUID] = Field(default_factory=list) 

728 

729 

730class FlowRunOrchestrationResult(PrefectBaseModel): 

731 """Per-run result for bulk state operations.""" 

732 

733 flow_run_id: UUID 

734 status: SetStateStatus 

735 state: Optional[schemas.states.State] = None 

736 details: StateResponseDetails 

737 

738 

739class FlowRunBulkSetStateResponse(PrefectBaseModel): 

740 """Response from bulk set state operation.""" 

741 

742 results: List[FlowRunOrchestrationResult] = Field(default_factory=list) 

743 

744 

745class FlowRunCreateResult(PrefectBaseModel): 

746 """Per-run result for bulk create operations.""" 

747 

748 flow_run_id: Optional[UUID] = None 

749 status: Literal["CREATED", "FAILED"] 

750 error: Optional[str] = None 

751 

752 

753class FlowRunBulkCreateResponse(PrefectBaseModel): 

754 """Response from bulk flow run creation.""" 

755 

756 results: List[FlowRunCreateResult] = Field(default_factory=list) 

757 

758 

759class FlowRunSlotSummary(PrefectBaseModel): 

760 """Summary of a flow run occupying a concurrency slot.""" 

761 

762 id: UUID 

763 name: str 

764 state_type: Optional[schemas.states.StateType] = None 

765 state_name: Optional[str] = None 

766 start_time: Optional[DateTime] = None 

767 state_timestamp: Optional[DateTime] = None 

768 time_in_current_state: Optional[timedelta] = None 

769 

770 

771class WorkQueueConcurrencyStatusDetail(PrefectBaseModel): 

772 """Per-queue concurrency status with flow run details.""" 

773 

774 queue_id: UUID 

775 queue_name: str 

776 active_slots: int 

777 concurrency_limit: Optional[int] = None 

778 flow_runs: List[FlowRunSlotSummary] = Field(default_factory=list) 

779 flow_run_count: Optional[int] = Field( 

780 default=None, 

781 description="Total flow run count for this queue (may differ from len(flow_runs) when flow_run_limit is applied).", 

782 ) 

783 

784 

785class WorkPoolConcurrencyStatus(PrefectBaseModel): 

786 """Paginated pool-level concurrency status with per-queue breakdown.""" 

787 

788 active_slots: int 

789 concurrency_limit: Optional[int] = None 

790 queues: List[WorkQueueConcurrencyStatusDetail] = Field(default_factory=list) 

791 count: int = Field(description="Total number of queues.") 

792 limit: int = Field(description="Page size.") 

793 pages: int = Field(description="Total number of pages.") 

794 page: int = Field(description="Current page number (1-indexed).") 

795 

796 

797class WorkQueueConcurrencyStatus(PrefectBaseModel): 

798 """Paginated queue-level concurrency status with flow run details.""" 

799 

800 active_slots: int 

801 concurrency_limit: Optional[int] = None 

802 flow_runs: List[FlowRunSlotSummary] = Field(default_factory=list) 

803 count: int = Field(description="Total number of slot-holding flow runs.") 

804 limit: int = Field(description="Page size.") 

805 pages: int = Field(description="Total number of pages.") 

806 page: int = Field(description="Current page number (1-indexed).")