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

1068 statements  

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

1""" 

2Schemas that define Prefect REST API filtering operations. 

3 

4Each filter schema includes logic for transforming itself into a SQL `where` clause. 

5""" 

6 

7from collections.abc import Iterable, Sequence 

8from typing import TYPE_CHECKING, ClassVar, Optional 

9from uuid import UUID 

10 

11from pydantic import ConfigDict, Field 

12from sqlalchemy.dialects import postgresql 

13from sqlalchemy.sql.functions import coalesce 

14 

15import prefect.server.schemas as schemas 

16from prefect.server.utilities.schemas.bases import PrefectBaseModel 

17from prefect.server.utilities.text_search_parser import ( 

18 parse_text_search_query, 

19) 

20from prefect.types import DateTime 

21from prefect.utilities.collections import AutoEnum 

22from prefect.utilities.importtools import lazy_import 

23 

24if TYPE_CHECKING: 24 ↛ 25line 24 didn't jump to line 25 because the condition on line 24 was never true

25 import sqlalchemy as sa 

26 

27 from prefect.server.database import PrefectDBInterface 

28 from prefect.server.schemas.core import Log 

29else: 

30 sa = lazy_import("sqlalchemy") 

31 

32# TODO: Consider moving the `as_sql_filter` functions out of here since they are a 

33# database model level function and do not properly separate concerns when 

34# present in the schemas module 

35 

36 

37def _as_array(elems: Sequence[str]) -> sa.ColumnElement[Sequence[str]]: 

38 return sa.cast(postgresql.array(elems), type_=postgresql.ARRAY(sa.String())) 

39 

40 

41class Operator(AutoEnum): 

42 """Operators for combining filter criteria.""" 

43 

44 and_ = AutoEnum.auto() 

45 or_ = AutoEnum.auto() 

46 

47 

48class PrefectFilterBaseModel(PrefectBaseModel): 

49 """Base model for Prefect filters""" 

50 

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

52 

53 def as_sql_filter(self) -> sa.ColumnElement[bool]: 

54 """Generate SQL filter from provided filter parameters. If no filters parameters are available, return a TRUE filter.""" 

55 from prefect.server.database.dependencies import provide_database_interface 

56 

57 db = provide_database_interface() 

58 filters = self._get_filter_list(db) 

59 if not filters: 

60 return sa.true() 

61 return sa.and_(*filters) 

62 

63 def _get_filter_list( 

64 self, db: "PrefectDBInterface" 

65 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

66 """Return a list of boolean filter statements based on filter parameters""" 

67 raise NotImplementedError("_get_filter_list must be implemented") 

68 

69 

70class PrefectOperatorFilterBaseModel(PrefectFilterBaseModel): 

71 """Base model for Prefect filters that combines criteria with a user-provided operator""" 

72 

73 operator: Operator = Field( 

74 default=Operator.and_, 

75 description="Operator for combining filter criteria. Defaults to 'and_'.", 

76 ) 

77 

78 def as_sql_filter(self) -> sa.ColumnElement[bool]: 

79 from prefect.server.database.dependencies import provide_database_interface 

80 

81 db = provide_database_interface() 

82 filters = self._get_filter_list(db) 

83 if not filters: 

84 return sa.true() 

85 return sa.and_(*filters) if self.operator == Operator.and_ else sa.or_(*filters) 

86 

87 

88class FlowFilterId(PrefectFilterBaseModel): 

89 """Filter by `Flow.id`.""" 

90 

91 any_: Optional[list[UUID]] = Field( 

92 default=None, description="A list of flow ids to include" 

93 ) 

94 not_any_: Optional[list[UUID]] = Field( 

95 default=None, description="A list of flow ids to exclude" 

96 ) 

97 

98 def _get_filter_list( 

99 self, db: "PrefectDBInterface" 

100 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

101 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

102 if self.any_ is not None: 

103 filters.append(db.Flow.id.in_(self.any_)) 

104 if self.not_any_: 

105 filters.append(db.Flow.id.not_in(self.not_any_)) 

106 return filters 

107 

108 

109class FlowFilterDeployment(PrefectOperatorFilterBaseModel): 

110 """Filter by flows by deployment""" 

111 

112 is_null_: Optional[bool] = Field( 

113 default=None, 

114 description="If true, only include flows without deployments", 

115 ) 

116 

117 def _get_filter_list( 

118 self, db: "PrefectDBInterface" 

119 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

120 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

121 

122 if self.is_null_ is not None: 

123 deployments_subquery = ( 

124 sa.select(db.Deployment.flow_id).distinct().subquery() 

125 ) 

126 

127 if self.is_null_: 

128 filters.append( 

129 db.Flow.id.not_in(sa.select(deployments_subquery.c.flow_id)) 

130 ) 

131 else: 

132 filters.append( 

133 db.Flow.id.in_(sa.select(deployments_subquery.c.flow_id)) 

134 ) 

135 

136 return filters 

137 

138 

139class FlowFilterName(PrefectFilterBaseModel): 

140 """Filter by `Flow.name`.""" 

141 

142 any_: Optional[list[str]] = Field( 

143 default=None, 

144 description="A list of flow names to include", 

145 examples=[["my-flow-1", "my-flow-2"]], 

146 ) 

147 

148 like_: Optional[str] = Field( 

149 default=None, 

150 description=( 

151 "A case-insensitive partial match. For example, " 

152 " passing 'marvin' will match " 

153 "'marvin', 'sad-Marvin', and 'marvin-robot'." 

154 ), 

155 examples=["marvin"], 

156 ) 

157 

158 def _get_filter_list( 

159 self, db: "PrefectDBInterface" 

160 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

161 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

162 if self.any_ is not None: 

163 filters.append(db.Flow.name.in_(self.any_)) 

164 if self.like_: 

165 filters.append(db.Flow.name.ilike(f"%{self.like_}%")) 

166 return filters 

167 

168 

169class FlowFilterTags(PrefectOperatorFilterBaseModel): 

170 """Filter by `Flow.tags`.""" 

171 

172 all_: Optional[list[str]] = Field( 

173 default=None, 

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

175 description=( 

176 "A list of tags. Flows will be returned only if their tags are a superset" 

177 " of the list" 

178 ), 

179 ) 

180 is_null_: Optional[bool] = Field( 

181 default=None, description="If true, only include flows without tags" 

182 ) 

183 

184 def _get_filter_list( 

185 self, db: "PrefectDBInterface" 

186 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

187 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

188 if self.all_: 

189 filters.append(db.Flow.tags.has_all(_as_array(self.all_))) 

190 if self.is_null_ is not None: 

191 filters.append(db.Flow.tags == [] if self.is_null_ else db.Flow.tags != []) 

192 return filters 

193 

194 

195class FlowFilter(PrefectOperatorFilterBaseModel): 

196 """Filter for flows. Only flows matching all criteria will be returned.""" 

197 

198 id: Optional[FlowFilterId] = Field( 

199 default=None, description="Filter criteria for `Flow.id`" 

200 ) 

201 deployment: Optional[FlowFilterDeployment] = Field( 

202 default=None, description="Filter criteria for Flow deployments" 

203 ) 

204 name: Optional[FlowFilterName] = Field( 

205 default=None, description="Filter criteria for `Flow.name`" 

206 ) 

207 tags: Optional[FlowFilterTags] = Field( 

208 default=None, description="Filter criteria for `Flow.tags`" 

209 ) 

210 

211 def _get_filter_list( 

212 self, db: "PrefectDBInterface" 

213 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

214 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

215 

216 if self.id is not None: 

217 filters.append(self.id.as_sql_filter()) 

218 if self.deployment is not None: 

219 filters.append(self.deployment.as_sql_filter()) 

220 if self.name is not None: 

221 filters.append(self.name.as_sql_filter()) 

222 if self.tags is not None: 

223 filters.append(self.tags.as_sql_filter()) 

224 

225 return filters 

226 

227 

228class FlowRunFilterId(PrefectFilterBaseModel): 

229 """Filter by `FlowRun.id`.""" 

230 

231 any_: Optional[list[UUID]] = Field( 

232 default=None, description="A list of flow run ids to include" 

233 ) 

234 not_any_: Optional[list[UUID]] = Field( 

235 default=None, description="A list of flow run ids to exclude" 

236 ) 

237 

238 def _get_filter_list( 

239 self, db: "PrefectDBInterface" 

240 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

241 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

242 if self.any_ is not None: 

243 filters.append(db.FlowRun.id.in_(self.any_)) 

244 if self.not_any_: 

245 filters.append(db.FlowRun.id.not_in(self.not_any_)) 

246 return filters 

247 

248 

249class FlowRunFilterName(PrefectFilterBaseModel): 

250 """Filter by `FlowRun.name`.""" 

251 

252 any_: Optional[list[str]] = Field( 

253 default=None, 

254 description="A list of flow run names to include", 

255 examples=[["my-flow-run-1", "my-flow-run-2"]], 

256 ) 

257 

258 like_: Optional[str] = Field( 

259 default=None, 

260 description=( 

261 "A case-insensitive partial match. For example, " 

262 " passing 'marvin' will match " 

263 "'marvin', 'sad-Marvin', and 'marvin-robot'." 

264 ), 

265 examples=["marvin"], 

266 ) 

267 

268 def _get_filter_list( 

269 self, db: "PrefectDBInterface" 

270 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

271 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

272 if self.any_ is not None: 

273 filters.append(db.FlowRun.name.in_(self.any_)) 

274 if self.like_: 

275 filters.append(db.FlowRun.name.ilike(f"%{self.like_}%")) 

276 return filters 

277 

278 

279class FlowRunFilterTags(PrefectOperatorFilterBaseModel): 

280 """Filter by `FlowRun.tags`.""" 

281 

282 all_: Optional[list[str]] = Field( 

283 default=None, 

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

285 description=( 

286 "A list of tags. Flow runs will be returned only if their tags are a" 

287 " superset of the list" 

288 ), 

289 ) 

290 

291 any_: Optional[list[str]] = Field( 

292 default=None, 

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

294 description="A list of tags to include", 

295 ) 

296 

297 is_null_: Optional[bool] = Field( 

298 default=None, description="If true, only include flow runs without tags" 

299 ) 

300 

301 def _get_filter_list( 

302 self, db: "PrefectDBInterface" 

303 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

304 def as_array(elems: Sequence[str]) -> sa.ColumnElement[Sequence[str]]: 

305 return sa.cast(postgresql.array(elems), type_=postgresql.ARRAY(sa.String())) 

306 

307 filters: list[sa.ColumnElement[bool]] = [] 

308 if self.all_: 

309 filters.append(db.FlowRun.tags.has_all(as_array(self.all_))) 

310 if self.any_ is not None: 

311 filters.append(db.FlowRun.tags.has_any(as_array(self.any_))) 

312 if self.is_null_ is not None: 

313 filters.append( 

314 db.FlowRun.tags == [] if self.is_null_ else db.FlowRun.tags != [] 

315 ) 

316 return filters 

317 

318 

319class FlowRunFilterDeploymentId(PrefectOperatorFilterBaseModel): 

320 """Filter by `FlowRun.deployment_id`.""" 

321 

322 any_: Optional[list[UUID]] = Field( 

323 default=None, description="A list of flow run deployment ids to include" 

324 ) 

325 is_null_: Optional[bool] = Field( 

326 default=None, 

327 description="If true, only include flow runs without deployment ids", 

328 ) 

329 

330 def _get_filter_list( 

331 self, db: "PrefectDBInterface" 

332 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

333 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

334 if self.any_ is not None: 

335 filters.append(db.FlowRun.deployment_id.in_(self.any_)) 

336 if self.is_null_ is not None: 

337 filters.append( 

338 db.FlowRun.deployment_id.is_(None) 

339 if self.is_null_ 

340 else db.FlowRun.deployment_id.is_not(None) 

341 ) 

342 return filters 

343 

344 

345class FlowRunFilterWorkQueueName(PrefectOperatorFilterBaseModel): 

346 """Filter by `FlowRun.work_queue_name`.""" 

347 

348 any_: Optional[list[str]] = Field( 

349 default=None, 

350 description="A list of work queue names to include", 

351 examples=[["work_queue_1", "work_queue_2"]], 

352 ) 

353 is_null_: Optional[bool] = Field( 

354 default=None, 

355 description="If true, only include flow runs without work queue names", 

356 ) 

357 

358 def _get_filter_list( 

359 self, db: "PrefectDBInterface" 

360 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

361 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

362 if self.any_ is not None: 

363 filters.append(db.FlowRun.work_queue_name.in_(self.any_)) 

364 if self.is_null_ is not None: 

365 filters.append( 

366 db.FlowRun.work_queue_name.is_(None) 

367 if self.is_null_ 

368 else db.FlowRun.work_queue_name.is_not(None) 

369 ) 

370 return filters 

371 

372 

373class FlowRunFilterStateType(PrefectFilterBaseModel): 

374 """Filter by `FlowRun.state_type`.""" 

375 

376 any_: Optional[list[schemas.states.StateType]] = Field( 

377 default=None, description="A list of flow run state types to include" 

378 ) 

379 not_any_: Optional[list[schemas.states.StateType]] = Field( 

380 default=None, description="A list of flow run state types to exclude" 

381 ) 

382 

383 def _get_filter_list( 

384 self, db: "PrefectDBInterface" 

385 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

386 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

387 if self.any_ is not None: 

388 filters.append(db.FlowRun.state_type.in_(self.any_)) 

389 if self.not_any_: 

390 filters.append(db.FlowRun.state_type.not_in(self.not_any_)) 

391 return filters 

392 

393 

394class FlowRunFilterStateName(PrefectFilterBaseModel): 

395 """Filter by `FlowRun.state_name`.""" 

396 

397 any_: Optional[list[str]] = Field( 

398 default=None, description="A list of flow run state names to include" 

399 ) 

400 not_any_: Optional[list[str]] = Field( 

401 default=None, description="A list of flow run state names to exclude" 

402 ) 

403 

404 def _get_filter_list( 

405 self, db: "PrefectDBInterface" 

406 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

407 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

408 if self.any_ is not None: 

409 filters.append(db.FlowRun.state_name.in_(self.any_)) 

410 if self.not_any_: 

411 filters.append(db.FlowRun.state_name.not_in(self.not_any_)) 

412 return filters 

413 

414 

415class FlowRunFilterState(PrefectOperatorFilterBaseModel): 

416 """Filter by `FlowRun.state_type` and `FlowRun.state_name`.""" 

417 

418 type: Optional[FlowRunFilterStateType] = Field( 

419 default=None, description="Filter criteria for `FlowRun.state_type`" 

420 ) 

421 name: Optional[FlowRunFilterStateName] = Field( 

422 default=None, description="Filter criteria for `FlowRun.state_name`" 

423 ) 

424 

425 def _get_filter_list( 

426 self, db: "PrefectDBInterface" 

427 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

428 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

429 if self.type is not None: 

430 filter = self.type.as_sql_filter() 

431 if isinstance(filter, sa.BinaryExpression): 

432 filters.append(filter) 

433 if self.name is not None: 

434 filter = self.name.as_sql_filter() 

435 if isinstance(filter, sa.BinaryExpression): 

436 filters.append(filter) 

437 return filters 

438 

439 

440class FlowRunFilterFlowVersion(PrefectFilterBaseModel): 

441 """Filter by `FlowRun.flow_version`.""" 

442 

443 any_: Optional[list[str]] = Field( 

444 default=None, description="A list of flow run flow_versions to include" 

445 ) 

446 

447 def _get_filter_list( 

448 self, db: "PrefectDBInterface" 

449 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

450 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

451 if self.any_ is not None: 

452 filters.append(db.FlowRun.flow_version.in_(self.any_)) 

453 return filters 

454 

455 

456class FlowRunFilterStartTime(PrefectFilterBaseModel): 

457 """Filter by `FlowRun.start_time`.""" 

458 

459 before_: Optional[DateTime] = Field( 

460 default=None, 

461 description="Only include flow runs starting at or before this time", 

462 ) 

463 after_: Optional[DateTime] = Field( 

464 default=None, 

465 description="Only include flow runs starting at or after this time", 

466 ) 

467 is_null_: Optional[bool] = Field( 

468 default=None, description="If true, only return flow runs without a start time" 

469 ) 

470 

471 def _get_filter_list( 

472 self, db: "PrefectDBInterface" 

473 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

474 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

475 if self.before_ is not None: 

476 filters.append( 

477 coalesce(db.FlowRun.start_time, db.FlowRun.expected_start_time) 

478 <= self.before_ 

479 ) 

480 if self.after_ is not None: 

481 filters.append( 

482 coalesce(db.FlowRun.start_time, db.FlowRun.expected_start_time) 

483 >= self.after_ 

484 ) 

485 if self.is_null_ is not None: 

486 filters.append( 

487 db.FlowRun.start_time.is_(None) 

488 if self.is_null_ 

489 else db.FlowRun.start_time.is_not(None) 

490 ) 

491 return filters 

492 

493 

494class FlowRunFilterEndTime(PrefectFilterBaseModel): 

495 """Filter by `FlowRun.end_time`.""" 

496 

497 before_: Optional[DateTime] = Field( 

498 default=None, 

499 description="Only include flow runs ending at or before this time", 

500 ) 

501 after_: Optional[DateTime] = Field( 

502 default=None, 

503 description="Only include flow runs ending at or after this time", 

504 ) 

505 is_null_: Optional[bool] = Field( 

506 default=None, description="If true, only return flow runs without an end time" 

507 ) 

508 

509 def _get_filter_list( 

510 self, db: "PrefectDBInterface" 

511 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

512 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

513 if self.before_ is not None: 

514 filters.append(db.FlowRun.end_time <= self.before_) 

515 if self.after_ is not None: 

516 filters.append(db.FlowRun.end_time >= self.after_) 

517 if self.is_null_ is not None: 

518 filters.append( 

519 db.FlowRun.end_time.is_(None) 

520 if self.is_null_ 

521 else db.FlowRun.end_time.is_not(None) 

522 ) 

523 return filters 

524 

525 

526class FlowRunFilterExpectedStartTime(PrefectFilterBaseModel): 

527 """Filter by `FlowRun.expected_start_time`.""" 

528 

529 before_: Optional[DateTime] = Field( 

530 default=None, 

531 description="Only include flow runs scheduled to start at or before this time", 

532 ) 

533 after_: Optional[DateTime] = Field( 

534 default=None, 

535 description="Only include flow runs scheduled to start at or after this time", 

536 ) 

537 

538 def _get_filter_list( 

539 self, db: "PrefectDBInterface" 

540 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

541 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

542 if self.before_ is not None: 

543 filters.append(db.FlowRun.expected_start_time <= self.before_) 

544 if self.after_ is not None: 

545 filters.append(db.FlowRun.expected_start_time >= self.after_) 

546 return filters 

547 

548 

549class FlowRunFilterNextScheduledStartTime(PrefectFilterBaseModel): 

550 """Filter by `FlowRun.next_scheduled_start_time`.""" 

551 

552 before_: Optional[DateTime] = Field( 

553 default=None, 

554 description=( 

555 "Only include flow runs with a next_scheduled_start_time or before this" 

556 " time" 

557 ), 

558 ) 

559 after_: Optional[DateTime] = Field( 

560 default=None, 

561 description=( 

562 "Only include flow runs with a next_scheduled_start_time at or after this" 

563 " time" 

564 ), 

565 ) 

566 

567 def _get_filter_list( 

568 self, db: "PrefectDBInterface" 

569 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

570 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

571 if self.before_ is not None: 

572 filters.append(db.FlowRun.next_scheduled_start_time <= self.before_) 

573 if self.after_ is not None: 

574 filters.append(db.FlowRun.next_scheduled_start_time >= self.after_) 

575 return filters 

576 

577 

578class FlowRunFilterParentFlowRunId(PrefectOperatorFilterBaseModel): 

579 """Filter for subflows of a given flow run""" 

580 

581 any_: Optional[list[UUID]] = Field( 

582 default=None, description="A list of parent flow run ids to include" 

583 ) 

584 

585 def _get_filter_list( 

586 self, db: "PrefectDBInterface" 

587 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

588 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

589 if self.any_ is not None: 

590 filters.append( 

591 db.FlowRun.id.in_( 

592 sa.select(db.FlowRun.id) 

593 .join( 

594 db.TaskRun, 

595 sa.and_( 

596 db.TaskRun.id == db.FlowRun.parent_task_run_id, 

597 ), 

598 ) 

599 .where(db.TaskRun.flow_run_id.in_(self.any_)) 

600 ) 

601 ) 

602 return filters 

603 

604 

605class FlowRunFilterParentTaskRunId(PrefectOperatorFilterBaseModel): 

606 """Filter by `FlowRun.parent_task_run_id`.""" 

607 

608 any_: Optional[list[UUID]] = Field( 

609 default=None, description="A list of flow run parent_task_run_ids to include" 

610 ) 

611 is_null_: Optional[bool] = Field( 

612 default=None, 

613 description="If true, only include flow runs without parent_task_run_id", 

614 ) 

615 

616 def _get_filter_list( 

617 self, db: "PrefectDBInterface" 

618 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

619 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

620 if self.any_ is not None: 

621 filters.append(db.FlowRun.parent_task_run_id.in_(self.any_)) 

622 if self.is_null_ is not None: 

623 filters.append( 

624 db.FlowRun.parent_task_run_id.is_(None) 

625 if self.is_null_ 

626 else db.FlowRun.parent_task_run_id.is_not(None) 

627 ) 

628 return filters 

629 

630 

631class FlowRunFilterIdempotencyKey(PrefectFilterBaseModel): 

632 """Filter by FlowRun.idempotency_key.""" 

633 

634 any_: Optional[list[str]] = Field( 

635 default=None, description="A list of flow run idempotency keys to include" 

636 ) 

637 not_any_: Optional[list[str]] = Field( 

638 default=None, description="A list of flow run idempotency keys to exclude" 

639 ) 

640 

641 def _get_filter_list( 

642 self, db: "PrefectDBInterface" 

643 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

644 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

645 if self.any_ is not None: 

646 filters.append(db.FlowRun.idempotency_key.in_(self.any_)) 

647 if self.not_any_: 

648 filters.append(db.FlowRun.idempotency_key.not_in(self.not_any_)) 

649 return filters 

650 

651 

652class FlowRunFilterCreatedBy(PrefectOperatorFilterBaseModel): 

653 """Filter by `FlowRun.created_by`.""" 

654 

655 id_: Optional[list[UUID]] = Field( 

656 default=None, 

657 description="A list of creator IDs to include", 

658 ) 

659 type_: Optional[list[str]] = Field( 

660 default=None, 

661 description=( 

662 "A list of creator types to include. For example, 'DEPLOYMENT' for " 

663 "scheduled runs or 'AUTOMATION' for runs triggered by automations." 

664 ), 

665 examples=[["DEPLOYMENT", "AUTOMATION"]], 

666 ) 

667 is_null_: Optional[bool] = Field( 

668 default=None, 

669 description="If true, only include flow runs without a creator", 

670 ) 

671 

672 def _get_filter_list( 

673 self, db: "PrefectDBInterface" 

674 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

675 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

676 if self.id_ is not None: 

677 # JSON stores UUIDs as strings, use astext for text extraction 

678 id_strings = [str(id_val) for id_val in self.id_] 

679 filters.append(db.FlowRun.created_by["id"].astext.in_(id_strings)) 

680 if self.type_ is not None: 

681 filters.append(db.FlowRun.created_by["type"].astext.in_(self.type_)) 

682 if self.is_null_ is not None: 

683 filters.append( 

684 db.FlowRun.created_by.is_(None) 

685 if self.is_null_ 

686 else db.FlowRun.created_by.is_not(None) 

687 ) 

688 return filters 

689 

690 

691class FlowRunFilter(PrefectOperatorFilterBaseModel): 

692 """Filter flow runs. Only flow runs matching all criteria will be returned""" 

693 

694 id: Optional[FlowRunFilterId] = Field( 

695 default=None, description="Filter criteria for `FlowRun.id`" 

696 ) 

697 name: Optional[FlowRunFilterName] = Field( 

698 default=None, description="Filter criteria for `FlowRun.name`" 

699 ) 

700 tags: Optional[FlowRunFilterTags] = Field( 

701 default=None, description="Filter criteria for `FlowRun.tags`" 

702 ) 

703 deployment_id: Optional[FlowRunFilterDeploymentId] = Field( 

704 default=None, description="Filter criteria for `FlowRun.deployment_id`" 

705 ) 

706 work_queue_name: Optional[FlowRunFilterWorkQueueName] = Field( 

707 default=None, description="Filter criteria for `FlowRun.work_queue_name" 

708 ) 

709 state: Optional[FlowRunFilterState] = Field( 

710 default=None, description="Filter criteria for `FlowRun.state`" 

711 ) 

712 flow_version: Optional[FlowRunFilterFlowVersion] = Field( 

713 default=None, description="Filter criteria for `FlowRun.flow_version`" 

714 ) 

715 start_time: Optional[FlowRunFilterStartTime] = Field( 

716 default=None, description="Filter criteria for `FlowRun.start_time`" 

717 ) 

718 end_time: Optional[FlowRunFilterEndTime] = Field( 

719 default=None, description="Filter criteria for `FlowRun.end_time`" 

720 ) 

721 expected_start_time: Optional[FlowRunFilterExpectedStartTime] = Field( 

722 default=None, description="Filter criteria for `FlowRun.expected_start_time`" 

723 ) 

724 next_scheduled_start_time: Optional[FlowRunFilterNextScheduledStartTime] = Field( 

725 default=None, 

726 description="Filter criteria for `FlowRun.next_scheduled_start_time`", 

727 ) 

728 parent_flow_run_id: Optional[FlowRunFilterParentFlowRunId] = Field( 

729 default=None, description="Filter criteria for subflows of the given flow runs" 

730 ) 

731 parent_task_run_id: Optional[FlowRunFilterParentTaskRunId] = Field( 

732 default=None, description="Filter criteria for `FlowRun.parent_task_run_id`" 

733 ) 

734 idempotency_key: Optional[FlowRunFilterIdempotencyKey] = Field( 

735 default=None, description="Filter criteria for `FlowRun.idempotency_key`" 

736 ) 

737 created_by: Optional[FlowRunFilterCreatedBy] = Field( 

738 default=None, description="Filter criteria for `FlowRun.created_by`" 

739 ) 

740 

741 def only_filters_on_id(self) -> bool: 

742 return bool( 

743 self.id is not None 

744 and (self.id.any_ and not self.id.not_any_) 

745 and self.name is None 

746 and self.tags is None 

747 and self.deployment_id is None 

748 and self.work_queue_name is None 

749 and self.state is None 

750 and self.flow_version is None 

751 and self.start_time is None 

752 and self.end_time is None 

753 and self.expected_start_time is None 

754 and self.next_scheduled_start_time is None 

755 and self.parent_flow_run_id is None 

756 and self.parent_task_run_id is None 

757 and self.idempotency_key is None 

758 and self.created_by is None 

759 ) 

760 

761 def _get_filter_list( 

762 self, db: "PrefectDBInterface" 

763 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

764 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

765 

766 if self.id is not None: 

767 filters.append(self.id.as_sql_filter()) 

768 if self.name is not None: 

769 filters.append(self.name.as_sql_filter()) 

770 if self.tags is not None: 

771 filters.append(self.tags.as_sql_filter()) 

772 if self.deployment_id is not None: 

773 filters.append(self.deployment_id.as_sql_filter()) 

774 if self.work_queue_name is not None: 

775 filters.append(self.work_queue_name.as_sql_filter()) 

776 if self.flow_version is not None: 

777 filters.append(self.flow_version.as_sql_filter()) 

778 if self.state is not None: 

779 filters.append(self.state.as_sql_filter()) 

780 if self.start_time is not None: 

781 filters.append(self.start_time.as_sql_filter()) 

782 if self.end_time is not None: 

783 filters.append(self.end_time.as_sql_filter()) 

784 if self.expected_start_time is not None: 

785 filters.append(self.expected_start_time.as_sql_filter()) 

786 if self.next_scheduled_start_time is not None: 

787 filters.append(self.next_scheduled_start_time.as_sql_filter()) 

788 if self.parent_flow_run_id is not None: 

789 filters.append(self.parent_flow_run_id.as_sql_filter()) 

790 if self.parent_task_run_id is not None: 

791 filters.append(self.parent_task_run_id.as_sql_filter()) 

792 if self.idempotency_key is not None: 

793 filters.append(self.idempotency_key.as_sql_filter()) 

794 if self.created_by is not None: 

795 filters.append(self.created_by.as_sql_filter()) 

796 

797 return filters 

798 

799 

800class TaskRunFilterFlowRunId(PrefectOperatorFilterBaseModel): 

801 """Filter by `TaskRun.flow_run_id`.""" 

802 

803 any_: Optional[list[UUID]] = Field( 

804 default=None, description="A list of task run flow run ids to include" 

805 ) 

806 

807 is_null_: Optional[bool] = Field( 

808 default=False, description="Filter for task runs with None as their flow run id" 

809 ) 

810 

811 def _get_filter_list( 

812 self, db: "PrefectDBInterface" 

813 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

814 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

815 if self.is_null_ is True: 

816 filters.append(db.TaskRun.flow_run_id.is_(None)) 

817 elif self.is_null_ is False and self.any_ is None: 

818 filters.append(db.TaskRun.flow_run_id.is_not(None)) 

819 else: 

820 if self.any_ is not None: 

821 filters.append(db.TaskRun.flow_run_id.in_(self.any_)) 

822 return filters 

823 

824 

825class TaskRunFilterId(PrefectFilterBaseModel): 

826 """Filter by `TaskRun.id`.""" 

827 

828 any_: Optional[list[UUID]] = Field( 

829 default=None, description="A list of task run ids to include" 

830 ) 

831 

832 def _get_filter_list( 

833 self, db: "PrefectDBInterface" 

834 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

835 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

836 if self.any_ is not None: 

837 filters.append(db.TaskRun.id.in_(self.any_)) 

838 return filters 

839 

840 

841class TaskRunFilterName(PrefectFilterBaseModel): 

842 """Filter by `TaskRun.name`.""" 

843 

844 any_: Optional[list[str]] = Field( 

845 default=None, 

846 description="A list of task run names to include", 

847 examples=[["my-task-run-1", "my-task-run-2"]], 

848 ) 

849 

850 like_: Optional[str] = Field( 

851 default=None, 

852 description=( 

853 "A case-insensitive partial match. For example, " 

854 " passing 'marvin' will match " 

855 "'marvin', 'sad-Marvin', and 'marvin-robot'." 

856 ), 

857 examples=["marvin"], 

858 ) 

859 

860 def _get_filter_list( 

861 self, db: "PrefectDBInterface" 

862 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

863 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

864 if self.any_ is not None: 

865 filters.append(db.TaskRun.name.in_(self.any_)) 

866 if self.like_: 

867 filters.append(db.TaskRun.name.ilike(f"%{self.like_}%")) 

868 return filters 

869 

870 

871class TaskRunFilterTags(PrefectOperatorFilterBaseModel): 

872 """Filter by `TaskRun.tags`.""" 

873 

874 all_: Optional[list[str]] = Field( 

875 default=None, 

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

877 description=( 

878 "A list of tags. Task runs will be returned only if their tags are a" 

879 " superset of the list" 

880 ), 

881 ) 

882 is_null_: Optional[bool] = Field( 

883 default=None, description="If true, only include task runs without tags" 

884 ) 

885 

886 def _get_filter_list( 

887 self, db: "PrefectDBInterface" 

888 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

889 filters: list[sa.ColumnElement[bool]] = [] 

890 if self.all_: 

891 filters.append(db.TaskRun.tags.has_all(_as_array(self.all_))) 

892 if self.is_null_ is not None: 

893 filters.append( 

894 db.TaskRun.tags == [] if self.is_null_ else db.TaskRun.tags != [] 

895 ) 

896 return filters 

897 

898 

899class TaskRunFilterStateType(PrefectFilterBaseModel): 

900 """Filter by `TaskRun.state_type`.""" 

901 

902 any_: Optional[list[schemas.states.StateType]] = Field( 

903 default=None, description="A list of task run state types to include" 

904 ) 

905 

906 def _get_filter_list( 

907 self, db: "PrefectDBInterface" 

908 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

909 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

910 if self.any_ is not None: 

911 filters.append(db.TaskRun.state_type.in_(self.any_)) 

912 return filters 

913 

914 

915class TaskRunFilterStateName(PrefectFilterBaseModel): 

916 """Filter by `TaskRun.state_name`.""" 

917 

918 any_: Optional[list[str]] = Field( 

919 default=None, description="A list of task run state names to include" 

920 ) 

921 

922 def _get_filter_list( 

923 self, db: "PrefectDBInterface" 

924 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

925 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

926 if self.any_ is not None: 

927 filters.append(db.TaskRun.state_name.in_(self.any_)) 

928 return filters 

929 

930 

931class TaskRunFilterState(PrefectOperatorFilterBaseModel): 

932 """Filter by `TaskRun.type` and `TaskRun.name`.""" 

933 

934 type: Optional[TaskRunFilterStateType] = Field( 

935 default=None, description="Filter criteria for `TaskRun.state_type`" 

936 ) 

937 name: Optional[TaskRunFilterStateName] = Field( 

938 default=None, description="Filter criteria for `TaskRun.state_name`" 

939 ) 

940 

941 def _get_filter_list( 

942 self, db: "PrefectDBInterface" 

943 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

944 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

945 if self.type is not None: 

946 filter = self.type.as_sql_filter() 

947 if isinstance(filter, sa.BinaryExpression): 

948 filters.append(filter) 

949 if self.name is not None: 

950 filter = self.name.as_sql_filter() 

951 if isinstance(filter, sa.BinaryExpression): 

952 filters.append(filter) 

953 return filters 

954 

955 

956class TaskRunFilterSubFlowRuns(PrefectFilterBaseModel): 

957 """Filter by `TaskRun.subflow_run`.""" 

958 

959 exists_: Optional[bool] = Field( 

960 default=None, 

961 description=( 

962 "If true, only include task runs that are subflow run parents; if false," 

963 " exclude parent task runs" 

964 ), 

965 ) 

966 

967 def _get_filter_list( 

968 self, db: "PrefectDBInterface" 

969 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

970 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

971 if self.exists_ is True: 

972 filters.append(db.TaskRun.subflow_run.has()) 

973 elif self.exists_ is False: 

974 filters.append(sa.not_(db.TaskRun.subflow_run.has())) 

975 return filters 

976 

977 

978class TaskRunFilterStartTime(PrefectFilterBaseModel): 

979 """Filter by `TaskRun.start_time`.""" 

980 

981 before_: Optional[DateTime] = Field( 

982 default=None, 

983 description="Only include task runs starting at or before this time", 

984 ) 

985 after_: Optional[DateTime] = Field( 

986 default=None, 

987 description="Only include task runs starting at or after this time", 

988 ) 

989 is_null_: Optional[bool] = Field( 

990 default=None, description="If true, only return task runs without a start time" 

991 ) 

992 

993 def _get_filter_list( 

994 self, db: "PrefectDBInterface" 

995 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

996 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

997 if self.before_ is not None: 

998 filters.append(db.TaskRun.start_time <= self.before_) 

999 if self.after_ is not None: 

1000 filters.append(db.TaskRun.start_time >= self.after_) 

1001 if self.is_null_ is not None: 

1002 filters.append( 

1003 db.TaskRun.start_time.is_(None) 

1004 if self.is_null_ 

1005 else db.TaskRun.start_time.is_not(None) 

1006 ) 

1007 return filters 

1008 

1009 

1010class TaskRunFilterEndTime(PrefectFilterBaseModel): 

1011 """Filter by `TaskRun.end_time`.""" 

1012 

1013 before_: Optional[DateTime] = Field( 

1014 default=None, 

1015 description="Only include task runs ending at or before this time", 

1016 ) 

1017 after_: Optional[DateTime] = Field( 

1018 default=None, 

1019 description="Only include task runs ending at or after this time", 

1020 ) 

1021 is_null_: Optional[bool] = Field( 

1022 default=None, description="If true, only return task runs without an end time" 

1023 ) 

1024 

1025 def _get_filter_list( 

1026 self, db: "PrefectDBInterface" 

1027 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

1028 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

1029 if self.before_ is not None: 

1030 filters.append(db.TaskRun.end_time <= self.before_) 

1031 if self.after_ is not None: 

1032 filters.append(db.TaskRun.end_time >= self.after_) 

1033 if self.is_null_ is not None: 

1034 filters.append( 

1035 db.TaskRun.end_time.is_(None) 

1036 if self.is_null_ 

1037 else db.TaskRun.end_time.is_not(None) 

1038 ) 

1039 return filters 

1040 

1041 

1042class TaskRunFilterExpectedStartTime(PrefectFilterBaseModel): 

1043 """Filter by `TaskRun.expected_start_time`.""" 

1044 

1045 before_: Optional[DateTime] = Field( 

1046 default=None, 

1047 description="Only include task runs expected to start at or before this time", 

1048 ) 

1049 after_: Optional[DateTime] = Field( 

1050 default=None, 

1051 description="Only include task runs expected to start at or after this time", 

1052 ) 

1053 

1054 def _get_filter_list( 

1055 self, db: "PrefectDBInterface" 

1056 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

1057 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

1058 if self.before_ is not None: 

1059 filters.append(db.TaskRun.expected_start_time <= self.before_) 

1060 if self.after_ is not None: 

1061 filters.append(db.TaskRun.expected_start_time >= self.after_) 

1062 return filters 

1063 

1064 

1065class TaskRunFilter(PrefectOperatorFilterBaseModel): 

1066 """Filter task runs. Only task runs matching all criteria will be returned""" 

1067 

1068 id: Optional[TaskRunFilterId] = Field( 

1069 default=None, description="Filter criteria for `TaskRun.id`" 

1070 ) 

1071 name: Optional[TaskRunFilterName] = Field( 

1072 default=None, description="Filter criteria for `TaskRun.name`" 

1073 ) 

1074 tags: Optional[TaskRunFilterTags] = Field( 

1075 default=None, description="Filter criteria for `TaskRun.tags`" 

1076 ) 

1077 state: Optional[TaskRunFilterState] = Field( 

1078 default=None, description="Filter criteria for `TaskRun.state`" 

1079 ) 

1080 start_time: Optional[TaskRunFilterStartTime] = Field( 

1081 default=None, description="Filter criteria for `TaskRun.start_time`" 

1082 ) 

1083 end_time: Optional[TaskRunFilterEndTime] = Field( 

1084 default=None, description="Filter criteria for `TaskRun.end_time`" 

1085 ) 

1086 expected_start_time: Optional[TaskRunFilterExpectedStartTime] = Field( 

1087 default=None, description="Filter criteria for `TaskRun.expected_start_time`" 

1088 ) 

1089 subflow_runs: Optional[TaskRunFilterSubFlowRuns] = Field( 

1090 default=None, description="Filter criteria for `TaskRun.subflow_run`" 

1091 ) 

1092 flow_run_id: Optional[TaskRunFilterFlowRunId] = Field( 

1093 default=None, description="Filter criteria for `TaskRun.flow_run_id`" 

1094 ) 

1095 

1096 def _get_filter_list( 

1097 self, db: "PrefectDBInterface" 

1098 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

1099 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

1100 

1101 if self.id is not None: 

1102 filters.append(self.id.as_sql_filter()) 

1103 if self.name is not None: 

1104 filters.append(self.name.as_sql_filter()) 

1105 if self.tags is not None: 

1106 filters.append(self.tags.as_sql_filter()) 

1107 if self.state is not None: 

1108 filters.append(self.state.as_sql_filter()) 

1109 if self.start_time is not None: 

1110 filters.append(self.start_time.as_sql_filter()) 

1111 if self.end_time is not None: 

1112 filters.append(self.end_time.as_sql_filter()) 

1113 if self.expected_start_time is not None: 

1114 filters.append(self.expected_start_time.as_sql_filter()) 

1115 if self.subflow_runs is not None: 

1116 filters.append(self.subflow_runs.as_sql_filter()) 

1117 if self.flow_run_id is not None: 

1118 filters.append(self.flow_run_id.as_sql_filter()) 

1119 

1120 return filters 

1121 

1122 

1123class DeploymentFilterId(PrefectFilterBaseModel): 

1124 """Filter by `Deployment.id`.""" 

1125 

1126 any_: Optional[list[UUID]] = Field( 

1127 default=None, description="A list of deployment ids to include" 

1128 ) 

1129 not_any_: Optional[list[UUID]] = Field( 

1130 default=None, description="A list of deployment ids to exclude" 

1131 ) 

1132 

1133 def _get_filter_list( 

1134 self, db: "PrefectDBInterface" 

1135 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

1136 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

1137 if self.any_ is not None: 

1138 filters.append(db.Deployment.id.in_(self.any_)) 

1139 if self.not_any_: 

1140 filters.append(db.Deployment.id.not_in(self.not_any_)) 

1141 return filters 

1142 

1143 

1144class DeploymentFilterName(PrefectFilterBaseModel): 

1145 """Filter by `Deployment.name`.""" 

1146 

1147 any_: Optional[list[str]] = Field( 

1148 default=None, 

1149 description="A list of deployment names to include", 

1150 examples=[["my-deployment-1", "my-deployment-2"]], 

1151 ) 

1152 

1153 like_: Optional[str] = Field( 

1154 default=None, 

1155 description=( 

1156 "A case-insensitive partial match. For example, " 

1157 " passing 'marvin' will match " 

1158 "'marvin', 'sad-Marvin', and 'marvin-robot'." 

1159 ), 

1160 examples=["marvin"], 

1161 ) 

1162 

1163 def _get_filter_list( 

1164 self, db: "PrefectDBInterface" 

1165 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

1166 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

1167 if self.any_ is not None: 

1168 filters.append(db.Deployment.name.in_(self.any_)) 

1169 if self.like_: 

1170 filters.append(db.Deployment.name.ilike(f"%{self.like_}%")) 

1171 return filters 

1172 

1173 

1174class DeploymentOrFlowNameFilter(PrefectFilterBaseModel): 

1175 """Filter by `Deployment.name` or `Flow.name` with a single input string for ilike filtering.""" 

1176 

1177 like_: Optional[str] = Field( 

1178 default=None, 

1179 description=( 

1180 "A case-insensitive partial match on deployment or flow names. For example, " 

1181 "passing 'example' might match deployments or flows with 'example' in their names." 

1182 ), 

1183 ) 

1184 

1185 def _get_filter_list( 

1186 self, db: "PrefectDBInterface" 

1187 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

1188 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

1189 if self.like_: 

1190 deployment_name_filter = db.Deployment.name.ilike(f"%{self.like_}%") 

1191 

1192 flow_name_filter = db.Deployment.flow.has( 

1193 db.Flow.name.ilike(f"%{self.like_}%") 

1194 ) 

1195 filters.append(sa.or_(deployment_name_filter, flow_name_filter)) 

1196 return filters 

1197 

1198 

1199class DeploymentFilterPaused(PrefectFilterBaseModel): 

1200 """Filter by `Deployment.paused`.""" 

1201 

1202 eq_: Optional[bool] = Field( 

1203 default=None, 

1204 description="Only returns where deployment is/is not paused", 

1205 ) 

1206 

1207 def _get_filter_list( 

1208 self, db: "PrefectDBInterface" 

1209 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

1210 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

1211 if self.eq_ is not None: 

1212 filters.append(db.Deployment.paused.is_(self.eq_)) 

1213 return filters 

1214 

1215 

1216class DeploymentFilterWorkQueueName(PrefectFilterBaseModel): 

1217 """Filter by `Deployment.work_queue_name`.""" 

1218 

1219 any_: Optional[list[str]] = Field( 

1220 default=None, 

1221 description="A list of work queue names to include", 

1222 examples=[["work_queue_1", "work_queue_2"]], 

1223 ) 

1224 

1225 def _get_filter_list( 

1226 self, db: "PrefectDBInterface" 

1227 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

1228 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

1229 if self.any_ is not None: 

1230 filters.append(db.Deployment.work_queue_name.in_(self.any_)) 

1231 return filters 

1232 

1233 

1234class DeploymentFilterConcurrencyLimit(PrefectFilterBaseModel): 

1235 """DEPRECATED: Prefer `Deployment.concurrency_limit_id` over `Deployment.concurrency_limit`.""" 

1236 

1237 ge_: Optional[int] = Field( 

1238 default=None, 

1239 description="Only include deployments with a concurrency limit greater than or equal to this value", 

1240 ) 

1241 

1242 le_: Optional[int] = Field( 

1243 default=None, 

1244 description="Only include deployments with a concurrency limit less than or equal to this value", 

1245 ) 

1246 is_null_: Optional[bool] = Field( 

1247 default=None, 

1248 description="If true, only include deployments without a concurrency limit", 

1249 ) 

1250 

1251 def _get_filter_list( 

1252 self, db: "PrefectDBInterface" 

1253 ) -> list[sa.ColumnElement[bool]]: 

1254 # This used to filter on an `int` column that was moved to a `ForeignKey` relationship 

1255 # This filter is now deprecated rather than support filtering on the new relationship 

1256 return [] 

1257 

1258 

1259class DeploymentFilterTags(PrefectOperatorFilterBaseModel): 

1260 """Filter by `Deployment.tags`.""" 

1261 

1262 all_: Optional[list[str]] = Field( 

1263 default=None, 

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

1265 description=( 

1266 "A list of tags. Deployments will be returned only if their tags are a" 

1267 " superset of the list" 

1268 ), 

1269 ) 

1270 any_: Optional[list[str]] = Field( 

1271 default=None, 

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

1273 description="A list of tags to include", 

1274 ) 

1275 

1276 is_null_: Optional[bool] = Field( 

1277 default=None, description="If true, only include deployments without tags" 

1278 ) 

1279 

1280 def _get_filter_list( 

1281 self, db: "PrefectDBInterface" 

1282 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

1283 from prefect.server.database import orm_models 

1284 

1285 filters: list[sa.ColumnElement[bool]] = [] 

1286 if self.all_: 

1287 filters.append(orm_models.Deployment.tags.has_all(_as_array(self.all_))) 

1288 if self.any_ is not None: 

1289 filters.append(orm_models.Deployment.tags.has_any(_as_array(self.any_))) 

1290 if self.is_null_ is not None: 

1291 filters.append( 

1292 db.Deployment.tags == [] if self.is_null_ else db.Deployment.tags != [] 

1293 ) 

1294 return filters 

1295 

1296 

1297class DeploymentFilter(PrefectOperatorFilterBaseModel): 

1298 """Filter for deployments. Only deployments matching all criteria will be returned.""" 

1299 

1300 id: Optional[DeploymentFilterId] = Field( 

1301 default=None, description="Filter criteria for `Deployment.id`" 

1302 ) 

1303 name: Optional[DeploymentFilterName] = Field( 

1304 default=None, description="Filter criteria for `Deployment.name`" 

1305 ) 

1306 flow_or_deployment_name: Optional[DeploymentOrFlowNameFilter] = Field( 

1307 default=None, description="Filter criteria for `Deployment.name` or `Flow.name`" 

1308 ) 

1309 paused: Optional[DeploymentFilterPaused] = Field( 

1310 default=None, description="Filter criteria for `Deployment.paused`" 

1311 ) 

1312 tags: Optional[DeploymentFilterTags] = Field( 

1313 default=None, description="Filter criteria for `Deployment.tags`" 

1314 ) 

1315 work_queue_name: Optional[DeploymentFilterWorkQueueName] = Field( 

1316 default=None, description="Filter criteria for `Deployment.work_queue_name`" 

1317 ) 

1318 concurrency_limit: Optional[DeploymentFilterConcurrencyLimit] = Field( 

1319 default=None, 

1320 description="DEPRECATED: Prefer `Deployment.concurrency_limit_id` over `Deployment.concurrency_limit`. If provided, will be ignored for backwards-compatibility. Will be removed after December 2024.", 

1321 deprecated=True, 

1322 ) 

1323 

1324 def _get_filter_list( 

1325 self, db: "PrefectDBInterface" 

1326 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

1327 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

1328 if self.id is not None: 

1329 filters.append(self.id.as_sql_filter()) 

1330 if self.name is not None: 

1331 filters.append(self.name.as_sql_filter()) 

1332 if self.flow_or_deployment_name is not None: 

1333 filters.append(self.flow_or_deployment_name.as_sql_filter()) 

1334 if self.paused is not None: 

1335 filters.append(self.paused.as_sql_filter()) 

1336 if self.tags is not None: 

1337 filters.append(self.tags.as_sql_filter()) 

1338 if self.work_queue_name is not None: 

1339 filters.append(self.work_queue_name.as_sql_filter()) 

1340 

1341 return filters 

1342 

1343 

1344class DeploymentScheduleFilterActive(PrefectFilterBaseModel): 

1345 """Filter by `DeploymentSchedule.active`.""" 

1346 

1347 eq_: Optional[bool] = Field( 

1348 default=None, 

1349 description="Only returns where deployment schedule is/is not active", 

1350 ) 

1351 

1352 def _get_filter_list( 

1353 self, db: "PrefectDBInterface" 

1354 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

1355 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

1356 if self.eq_ is not None: 

1357 filters.append(db.DeploymentSchedule.active.is_(self.eq_)) 

1358 return filters 

1359 

1360 

1361class DeploymentScheduleFilter(PrefectOperatorFilterBaseModel): 

1362 """Filter for deployments. Only deployments matching all criteria will be returned.""" 

1363 

1364 active: Optional[DeploymentScheduleFilterActive] = Field( 

1365 default=None, description="Filter criteria for `DeploymentSchedule.active`" 

1366 ) 

1367 

1368 def _get_filter_list( 

1369 self, db: "PrefectDBInterface" 

1370 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

1371 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

1372 

1373 if self.active is not None: 

1374 filters.append(self.active.as_sql_filter()) 

1375 

1376 return filters 

1377 

1378 

1379class LogFilterName(PrefectFilterBaseModel): 

1380 """Filter by `Log.name`.""" 

1381 

1382 any_: Optional[list[str]] = Field( 

1383 default=None, 

1384 description="A list of log names to include", 

1385 examples=[["prefect.logger.flow_runs", "prefect.logger.task_runs"]], 

1386 ) 

1387 

1388 def _get_filter_list( 

1389 self, db: "PrefectDBInterface" 

1390 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

1391 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

1392 if self.any_ is not None: 

1393 filters.append(db.Log.name.in_(self.any_)) 

1394 return filters 

1395 

1396 

1397class LogFilterLevel(PrefectFilterBaseModel): 

1398 """Filter by `Log.level`.""" 

1399 

1400 ge_: Optional[int] = Field( 

1401 default=None, 

1402 description="Include logs with a level greater than or equal to this level", 

1403 examples=[20], 

1404 ) 

1405 

1406 le_: Optional[int] = Field( 

1407 default=None, 

1408 description="Include logs with a level less than or equal to this level", 

1409 examples=[50], 

1410 ) 

1411 

1412 def _get_filter_list( 

1413 self, db: "PrefectDBInterface" 

1414 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

1415 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

1416 if self.ge_ is not None: 1416 ↛ 1418line 1416 didn't jump to line 1418 because the condition on line 1416 was always true

1417 filters.append(db.Log.level >= self.ge_) 

1418 if self.le_ is not None: 1418 ↛ 1420line 1418 didn't jump to line 1420 because the condition on line 1418 was always true

1419 filters.append(db.Log.level <= self.le_) 

1420 return filters 

1421 

1422 

1423class LogFilterTimestamp(PrefectFilterBaseModel): 

1424 """Filter by `Log.timestamp`.""" 

1425 

1426 before_: Optional[DateTime] = Field( 

1427 default=None, 

1428 description="Only include logs with a timestamp at or before this time", 

1429 ) 

1430 after_: Optional[DateTime] = Field( 

1431 default=None, 

1432 description="Only include logs with a timestamp at or after this time", 

1433 ) 

1434 

1435 def _get_filter_list( 

1436 self, db: "PrefectDBInterface" 

1437 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

1438 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

1439 if self.before_ is not None: 1439 ↛ 1440line 1439 didn't jump to line 1440 because the condition on line 1439 was never true

1440 filters.append(db.Log.timestamp <= self.before_) 

1441 if self.after_ is not None: 1441 ↛ 1443line 1441 didn't jump to line 1443 because the condition on line 1441 was always true

1442 filters.append(db.Log.timestamp >= self.after_) 

1443 return filters 

1444 

1445 

1446class LogFilterFlowRunId(PrefectFilterBaseModel): 

1447 """Filter by `Log.flow_run_id`.""" 

1448 

1449 any_: Optional[list[UUID]] = Field( 

1450 default=None, description="A list of flow run IDs to include" 

1451 ) 

1452 

1453 def _get_filter_list( 

1454 self, db: "PrefectDBInterface" 

1455 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

1456 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

1457 if self.any_ is not None: 1457 ↛ 1459line 1457 didn't jump to line 1459 because the condition on line 1457 was always true

1458 filters.append(db.Log.flow_run_id.in_(self.any_)) 

1459 return filters 

1460 

1461 

1462class LogFilterTaskRunId(PrefectFilterBaseModel): 

1463 """Filter by `Log.task_run_id`.""" 

1464 

1465 any_: Optional[list[UUID]] = Field( 

1466 default=None, description="A list of task run IDs to include" 

1467 ) 

1468 

1469 is_null_: Optional[bool] = Field( 

1470 default=None, 

1471 description="If true, only include logs without a task run id", 

1472 ) 

1473 

1474 def _get_filter_list( 

1475 self, db: "PrefectDBInterface" 

1476 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

1477 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

1478 if self.any_ is not None: 1478 ↛ 1480line 1478 didn't jump to line 1480 because the condition on line 1478 was always true

1479 filters.append(db.Log.task_run_id.in_(self.any_)) 

1480 if self.is_null_ is not None: 

1481 filters.append( 

1482 db.Log.task_run_id.is_(None) 

1483 if self.is_null_ 

1484 else db.Log.task_run_id.is_not(None) 

1485 ) 

1486 return filters 

1487 

1488 

1489class LogFilterTextSearch(PrefectFilterBaseModel): 

1490 """Filter by text search across log content.""" 

1491 

1492 query: str = Field( 

1493 description="Text search query string", 

1494 examples=[ 

1495 "error", 

1496 "error -debug", 

1497 '"connection timeout"', 

1498 "+required -excluded", 

1499 ], 

1500 max_length=200, 

1501 ) 

1502 

1503 def includes(self, log: "Log") -> bool: 

1504 """Check if this text filter includes the given log.""" 

1505 from prefect.server.schemas.core import Log 

1506 

1507 if not isinstance(log, Log): 

1508 raise TypeError(f"Expected Log object, got {type(log)}") 

1509 

1510 # Parse query into components 

1511 parsed = parse_text_search_query(self.query) 

1512 

1513 # Build searchable text from message and logger name 

1514 searchable_text = f"{log.message} {log.name}".lower() 

1515 

1516 # Check include terms (OR logic) 

1517 if parsed.include: 

1518 include_match = any( 

1519 term.lower() in searchable_text for term in parsed.include 

1520 ) 

1521 if not include_match: 

1522 return False 

1523 

1524 # Check exclude terms (NOT logic) 

1525 if parsed.exclude: 

1526 exclude_match = any( 

1527 term.lower() in searchable_text for term in parsed.exclude 

1528 ) 

1529 if exclude_match: 

1530 return False 

1531 

1532 # Check required terms (AND logic - future feature) 

1533 if parsed.required: 

1534 required_match = all( 

1535 term.lower() in searchable_text for term in parsed.required 

1536 ) 

1537 if not required_match: 

1538 return False 

1539 

1540 return True 

1541 

1542 def _get_filter_list( 

1543 self, db: "PrefectDBInterface" 

1544 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

1545 """Build SQLAlchemy WHERE clauses for text search""" 

1546 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

1547 

1548 if not self.query.strip(): 1548 ↛ 1549line 1548 didn't jump to line 1549 because the condition on line 1548 was never true

1549 return filters 

1550 

1551 parsed = parse_text_search_query(self.query) 

1552 

1553 # Build combined searchable text field (message + name) 

1554 searchable_field = sa.func.concat(db.Log.message, " ", db.Log.name) 

1555 

1556 # Handle include terms (OR logic) 

1557 if parsed.include: 

1558 include_conditions = [] 

1559 for term in parsed.include: 

1560 include_conditions.append( 

1561 sa.func.lower(searchable_field).contains(term.lower()) 

1562 ) 

1563 

1564 if include_conditions: 1564 ↛ 1568line 1564 didn't jump to line 1568 because the condition on line 1564 was always true

1565 filters.append(sa.or_(*include_conditions)) 

1566 

1567 # Handle exclude terms (NOT logic) 

1568 if parsed.exclude: 

1569 exclude_conditions = [] 

1570 for term in parsed.exclude: 

1571 exclude_conditions.append( 

1572 ~sa.func.lower(searchable_field).contains(term.lower()) 

1573 ) 

1574 

1575 if exclude_conditions: 1575 ↛ 1579line 1575 didn't jump to line 1579 because the condition on line 1575 was always true

1576 filters.append(sa.and_(*exclude_conditions)) 

1577 

1578 # Handle required terms (AND logic - future feature) 

1579 if parsed.required: 

1580 required_conditions = [] 

1581 for term in parsed.required: 

1582 required_conditions.append( 

1583 sa.func.lower(searchable_field).contains(term.lower()) 

1584 ) 

1585 

1586 if required_conditions: 1586 ↛ 1589line 1586 didn't jump to line 1589 because the condition on line 1586 was always true

1587 filters.append(sa.and_(*required_conditions)) 

1588 

1589 return filters 

1590 

1591 

1592class LogFilter(PrefectOperatorFilterBaseModel): 

1593 """Filter logs. Only logs matching all criteria will be returned""" 

1594 

1595 level: Optional[LogFilterLevel] = Field( 

1596 default=None, description="Filter criteria for `Log.level`" 

1597 ) 

1598 timestamp: Optional[LogFilterTimestamp] = Field( 

1599 default=None, description="Filter criteria for `Log.timestamp`" 

1600 ) 

1601 flow_run_id: Optional[LogFilterFlowRunId] = Field( 

1602 default=None, description="Filter criteria for `Log.flow_run_id`" 

1603 ) 

1604 task_run_id: Optional[LogFilterTaskRunId] = Field( 

1605 default=None, description="Filter criteria for `Log.task_run_id`" 

1606 ) 

1607 text: Optional[LogFilterTextSearch] = Field( 

1608 default=None, description="Filter criteria for text search across log content" 

1609 ) 

1610 

1611 def _get_filter_list( 

1612 self, db: "PrefectDBInterface" 

1613 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

1614 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

1615 

1616 if self.level is not None: 

1617 filters.append(self.level.as_sql_filter()) 

1618 if self.timestamp is not None: 

1619 filters.append(self.timestamp.as_sql_filter()) 

1620 if self.flow_run_id is not None: 

1621 filters.append(self.flow_run_id.as_sql_filter()) 

1622 if self.task_run_id is not None: 

1623 filters.append(self.task_run_id.as_sql_filter()) 

1624 if self.text is not None: 

1625 filters.extend(self.text._get_filter_list(db)) 

1626 

1627 return filters 

1628 

1629 

1630class FilterSet(PrefectBaseModel): 

1631 """A collection of filters for common objects""" 

1632 

1633 flows: FlowFilter = Field( 

1634 default_factory=FlowFilter, description="Filters that apply to flows" 

1635 ) 

1636 flow_runs: FlowRunFilter = Field( 

1637 default_factory=FlowRunFilter, description="Filters that apply to flow runs" 

1638 ) 

1639 task_runs: TaskRunFilter = Field( 

1640 default_factory=TaskRunFilter, description="Filters that apply to task runs" 

1641 ) 

1642 deployments: DeploymentFilter = Field( 

1643 default_factory=DeploymentFilter, 

1644 description="Filters that apply to deployments", 

1645 ) 

1646 

1647 

1648class BlockTypeFilterName(PrefectFilterBaseModel): 

1649 """Filter by `BlockType.name`""" 

1650 

1651 like_: Optional[str] = Field( 

1652 default=None, 

1653 description=( 

1654 "A case-insensitive partial match. For example, " 

1655 " passing 'marvin' will match " 

1656 "'marvin', 'sad-Marvin', and 'marvin-robot'." 

1657 ), 

1658 examples=["marvin"], 

1659 ) 

1660 

1661 def _get_filter_list( 

1662 self, db: "PrefectDBInterface" 

1663 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

1664 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

1665 if self.like_: 

1666 filters.append(db.BlockType.name.ilike(f"%{self.like_}%")) 

1667 return filters 

1668 

1669 

1670class BlockTypeFilterSlug(PrefectFilterBaseModel): 

1671 """Filter by `BlockType.slug`""" 

1672 

1673 any_: Optional[list[str]] = Field( 

1674 default=None, description="A list of slugs to match" 

1675 ) 

1676 

1677 def _get_filter_list( 

1678 self, db: "PrefectDBInterface" 

1679 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

1680 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

1681 if self.any_ is not None: 

1682 filters.append(db.BlockType.slug.in_(self.any_)) 

1683 

1684 return filters 

1685 

1686 

1687class BlockTypeFilter(PrefectFilterBaseModel): 

1688 """Filter BlockTypes""" 

1689 

1690 name: Optional[BlockTypeFilterName] = Field( 

1691 default=None, description="Filter criteria for `BlockType.name`" 

1692 ) 

1693 

1694 slug: Optional[BlockTypeFilterSlug] = Field( 

1695 default=None, description="Filter criteria for `BlockType.slug`" 

1696 ) 

1697 

1698 def _get_filter_list( 

1699 self, db: "PrefectDBInterface" 

1700 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

1701 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

1702 

1703 if self.name is not None: 

1704 filters.append(self.name.as_sql_filter()) 

1705 if self.slug is not None: 

1706 filters.append(self.slug.as_sql_filter()) 

1707 

1708 return filters 

1709 

1710 

1711class BlockSchemaFilterBlockTypeId(PrefectFilterBaseModel): 

1712 """Filter by `BlockSchema.block_type_id`.""" 

1713 

1714 any_: Optional[list[UUID]] = Field( 

1715 default=None, description="A list of block type ids to include" 

1716 ) 

1717 

1718 def _get_filter_list( 

1719 self, db: "PrefectDBInterface" 

1720 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

1721 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

1722 if self.any_ is not None: 1722 ↛ 1723line 1722 didn't jump to line 1723 because the condition on line 1722 was never true

1723 filters.append(db.BlockSchema.block_type_id.in_(self.any_)) 

1724 return filters 

1725 

1726 

1727class BlockSchemaFilterId(PrefectFilterBaseModel): 

1728 """Filter by BlockSchema.id""" 

1729 

1730 any_: Optional[list[UUID]] = Field( 

1731 default=None, description="A list of IDs to include" 

1732 ) 

1733 

1734 def _get_filter_list( 

1735 self, db: "PrefectDBInterface" 

1736 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

1737 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

1738 if self.any_ is not None: 

1739 filters.append(db.BlockSchema.id.in_(self.any_)) 

1740 return filters 

1741 

1742 

1743class BlockSchemaFilterCapabilities(PrefectFilterBaseModel): 

1744 """Filter by `BlockSchema.capabilities`""" 

1745 

1746 all_: Optional[list[str]] = Field( 

1747 default=None, 

1748 examples=[["write-storage", "read-storage"]], 

1749 description=( 

1750 "A list of block capabilities. Block entities will be returned only if an" 

1751 " associated block schema has a superset of the defined capabilities." 

1752 ), 

1753 ) 

1754 

1755 def _get_filter_list( 

1756 self, db: "PrefectDBInterface" 

1757 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

1758 filters: list[sa.ColumnElement[bool]] = [] 

1759 if self.all_: 

1760 filters.append(db.BlockSchema.capabilities.has_all(_as_array(self.all_))) 

1761 return filters 

1762 

1763 

1764class BlockSchemaFilterVersion(PrefectFilterBaseModel): 

1765 """Filter by `BlockSchema.capabilities`""" 

1766 

1767 any_: Optional[list[str]] = Field( 

1768 default=None, 

1769 examples=[["2.0.0", "2.1.0"]], 

1770 description="A list of block schema versions.", 

1771 ) 

1772 

1773 def _get_filter_list( 

1774 self, db: "PrefectDBInterface" 

1775 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

1776 filters: list[sa.ColumnElement[bool]] = [] 

1777 if self.any_ is not None: 

1778 filters.append(db.BlockSchema.version.in_(self.any_)) 

1779 return filters 

1780 

1781 

1782class BlockSchemaFilter(PrefectOperatorFilterBaseModel): 

1783 """Filter BlockSchemas""" 

1784 

1785 block_type_id: Optional[BlockSchemaFilterBlockTypeId] = Field( 

1786 default=None, description="Filter criteria for `BlockSchema.block_type_id`" 

1787 ) 

1788 block_capabilities: Optional[BlockSchemaFilterCapabilities] = Field( 

1789 default=None, description="Filter criteria for `BlockSchema.capabilities`" 

1790 ) 

1791 id: Optional[BlockSchemaFilterId] = Field( 

1792 default=None, description="Filter criteria for `BlockSchema.id`" 

1793 ) 

1794 version: Optional[BlockSchemaFilterVersion] = Field( 

1795 default=None, description="Filter criteria for `BlockSchema.version`" 

1796 ) 

1797 

1798 def _get_filter_list( 

1799 self, db: "PrefectDBInterface" 

1800 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

1801 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

1802 

1803 if self.block_type_id is not None: 

1804 filters.append(self.block_type_id.as_sql_filter()) 

1805 if self.block_capabilities is not None: 

1806 filters.append(self.block_capabilities.as_sql_filter()) 

1807 if self.id is not None: 

1808 filters.append(self.id.as_sql_filter()) 

1809 if self.version is not None: 

1810 filters.append(self.version.as_sql_filter()) 

1811 

1812 return filters 

1813 

1814 

1815class BlockDocumentFilterIsAnonymous(PrefectFilterBaseModel): 

1816 """Filter by `BlockDocument.is_anonymous`.""" 

1817 

1818 eq_: Optional[bool] = Field( 

1819 default=None, 

1820 description=( 

1821 "Filter block documents for only those that are or are not anonymous." 

1822 ), 

1823 ) 

1824 

1825 def _get_filter_list( 

1826 self, db: "PrefectDBInterface" 

1827 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

1828 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

1829 if self.eq_ is not None: 

1830 filters.append(db.BlockDocument.is_anonymous.is_(self.eq_)) 

1831 return filters 

1832 

1833 

1834class BlockDocumentFilterBlockTypeId(PrefectFilterBaseModel): 

1835 """Filter by `BlockDocument.block_type_id`.""" 

1836 

1837 any_: Optional[list[UUID]] = Field( 

1838 default=None, description="A list of block type ids to include" 

1839 ) 

1840 

1841 def _get_filter_list( 

1842 self, db: "PrefectDBInterface" 

1843 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

1844 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

1845 if self.any_ is not None: 

1846 filters.append(db.BlockDocument.block_type_id.in_(self.any_)) 

1847 return filters 

1848 

1849 

1850class BlockDocumentFilterId(PrefectFilterBaseModel): 

1851 """Filter by `BlockDocument.id`.""" 

1852 

1853 any_: Optional[list[UUID]] = Field( 

1854 default=None, description="A list of block ids to include" 

1855 ) 

1856 

1857 def _get_filter_list( 

1858 self, db: "PrefectDBInterface" 

1859 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

1860 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

1861 if self.any_ is not None: 

1862 filters.append(db.BlockDocument.id.in_(self.any_)) 

1863 return filters 

1864 

1865 

1866class BlockDocumentFilterName(PrefectFilterBaseModel): 

1867 """Filter by `BlockDocument.name`.""" 

1868 

1869 any_: Optional[list[str]] = Field( 

1870 default=None, description="A list of block names to include" 

1871 ) 

1872 like_: Optional[str] = Field( 

1873 default=None, 

1874 description=( 

1875 "A string to match block names against. This can include " 

1876 "SQL wildcard characters like `%` and `_`." 

1877 ), 

1878 examples=["my-block%"], 

1879 ) 

1880 

1881 def _get_filter_list( 

1882 self, db: "PrefectDBInterface" 

1883 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

1884 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

1885 if self.any_ is not None: 

1886 filters.append(db.BlockDocument.name.in_(self.any_)) 

1887 if self.like_: 

1888 filters.append(db.BlockDocument.name.ilike(f"%{self.like_}%")) 

1889 return filters 

1890 

1891 

1892class BlockDocumentFilter(PrefectOperatorFilterBaseModel): 

1893 """Filter BlockDocuments. Only BlockDocuments matching all criteria will be returned""" 

1894 

1895 id: Optional[BlockDocumentFilterId] = Field( 

1896 default=None, description="Filter criteria for `BlockDocument.id`" 

1897 ) 

1898 is_anonymous: Optional[BlockDocumentFilterIsAnonymous] = Field( 

1899 # default is to exclude anonymous blocks 

1900 BlockDocumentFilterIsAnonymous(eq_=False), 

1901 description=( 

1902 "Filter criteria for `BlockDocument.is_anonymous`. " 

1903 "Defaults to excluding anonymous blocks." 

1904 ), 

1905 ) 

1906 block_type_id: Optional[BlockDocumentFilterBlockTypeId] = Field( 

1907 default=None, description="Filter criteria for `BlockDocument.block_type_id`" 

1908 ) 

1909 name: Optional[BlockDocumentFilterName] = Field( 

1910 default=None, description="Filter criteria for `BlockDocument.name`" 

1911 ) 

1912 

1913 def _get_filter_list( 

1914 self, db: "PrefectDBInterface" 

1915 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

1916 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

1917 if self.id is not None: 

1918 filters.append(self.id.as_sql_filter()) 

1919 if self.is_anonymous is not None: 

1920 filters.append(self.is_anonymous.as_sql_filter()) 

1921 if self.block_type_id is not None: 

1922 filters.append(self.block_type_id.as_sql_filter()) 

1923 if self.name is not None: 

1924 filters.append(self.name.as_sql_filter()) 

1925 return filters 

1926 

1927 

1928class WorkQueueFilterId(PrefectFilterBaseModel): 

1929 """Filter by `WorkQueue.id`.""" 

1930 

1931 any_: Optional[list[UUID]] = Field( 

1932 default=None, 

1933 description="A list of work queue ids to include", 

1934 ) 

1935 

1936 def _get_filter_list( 

1937 self, db: "PrefectDBInterface" 

1938 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

1939 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

1940 if self.any_ is not None: 

1941 filters.append(db.WorkQueue.id.in_(self.any_)) 

1942 return filters 

1943 

1944 

1945class WorkQueueFilterName(PrefectFilterBaseModel): 

1946 """Filter by `WorkQueue.name`.""" 

1947 

1948 any_: Optional[list[str]] = Field( 

1949 default=None, 

1950 description="A list of work queue names to include", 

1951 examples=[["wq-1", "wq-2"]], 

1952 ) 

1953 

1954 startswith_: Optional[list[str]] = Field( 

1955 default=None, 

1956 description=( 

1957 "A list of case-insensitive starts-with matches. For example, " 

1958 " passing 'marvin' will match " 

1959 "'marvin', and 'Marvin-robot', but not 'sad-marvin'." 

1960 ), 

1961 examples=[["marvin", "Marvin-robot"]], 

1962 ) 

1963 

1964 def _get_filter_list( 

1965 self, db: "PrefectDBInterface" 

1966 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

1967 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

1968 if self.any_ is not None: 

1969 filters.append(db.WorkQueue.name.in_(self.any_)) 

1970 if self.startswith_ is not None: 

1971 filters.append( 

1972 sa.or_( 

1973 *[db.WorkQueue.name.ilike(f"{item}%") for item in self.startswith_] 

1974 ) 

1975 ) 

1976 return filters 

1977 

1978 

1979class WorkQueueFilter(PrefectOperatorFilterBaseModel): 

1980 """Filter work queues. Only work queues matching all criteria will be 

1981 returned""" 

1982 

1983 id: Optional[WorkQueueFilterId] = Field( 

1984 default=None, description="Filter criteria for `WorkQueue.id`" 

1985 ) 

1986 

1987 name: Optional[WorkQueueFilterName] = Field( 

1988 default=None, description="Filter criteria for `WorkQueue.name`" 

1989 ) 

1990 

1991 def _get_filter_list( 

1992 self, db: "PrefectDBInterface" 

1993 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

1994 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

1995 

1996 if self.id is not None: 

1997 filters.append(self.id.as_sql_filter()) 

1998 if self.name is not None: 

1999 filters.append(self.name.as_sql_filter()) 

2000 

2001 return filters 

2002 

2003 

2004class WorkPoolFilterId(PrefectFilterBaseModel): 

2005 """Filter by `WorkPool.id`.""" 

2006 

2007 any_: Optional[list[UUID]] = Field( 

2008 default=None, description="A list of work pool ids to include" 

2009 ) 

2010 

2011 def _get_filter_list( 

2012 self, db: "PrefectDBInterface" 

2013 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

2014 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

2015 if self.any_ is not None: 

2016 filters.append(db.WorkPool.id.in_(self.any_)) 

2017 return filters 

2018 

2019 

2020class WorkPoolFilterName(PrefectFilterBaseModel): 

2021 """Filter by `WorkPool.name`.""" 

2022 

2023 any_: Optional[list[str]] = Field( 

2024 default=None, description="A list of work pool names to include" 

2025 ) 

2026 

2027 def _get_filter_list( 

2028 self, db: "PrefectDBInterface" 

2029 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

2030 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

2031 if self.any_ is not None: 

2032 filters.append(db.WorkPool.name.in_(self.any_)) 

2033 return filters 

2034 

2035 

2036class WorkPoolFilterType(PrefectFilterBaseModel): 

2037 """Filter by `WorkPool.type`.""" 

2038 

2039 any_: Optional[list[str]] = Field( 

2040 default=None, description="A list of work pool types to include" 

2041 ) 

2042 

2043 def _get_filter_list( 

2044 self, db: "PrefectDBInterface" 

2045 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

2046 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

2047 if self.any_ is not None: 

2048 filters.append(db.WorkPool.type.in_(self.any_)) 

2049 return filters 

2050 

2051 

2052class WorkPoolFilter(PrefectOperatorFilterBaseModel): 

2053 """Filter work pools. Only work pools matching all criteria will be returned""" 

2054 

2055 id: Optional[WorkPoolFilterId] = Field( 

2056 default=None, description="Filter criteria for `WorkPool.id`" 

2057 ) 

2058 name: Optional[WorkPoolFilterName] = Field( 

2059 default=None, description="Filter criteria for `WorkPool.name`" 

2060 ) 

2061 type: Optional[WorkPoolFilterType] = Field( 

2062 default=None, description="Filter criteria for `WorkPool.type`" 

2063 ) 

2064 

2065 def _get_filter_list( 

2066 self, db: "PrefectDBInterface" 

2067 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

2068 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

2069 

2070 if self.id is not None: 

2071 filters.append(self.id.as_sql_filter()) 

2072 if self.name is not None: 

2073 filters.append(self.name.as_sql_filter()) 

2074 if self.type is not None: 

2075 filters.append(self.type.as_sql_filter()) 

2076 

2077 return filters 

2078 

2079 

2080class WorkerFilterWorkPoolId(PrefectFilterBaseModel): 

2081 """Filter by `Worker.worker_config_id`.""" 

2082 

2083 any_: Optional[list[UUID]] = Field( 

2084 default=None, description="A list of work pool ids to include" 

2085 ) 

2086 

2087 def _get_filter_list( 

2088 self, db: "PrefectDBInterface" 

2089 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

2090 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

2091 if self.any_ is not None: 

2092 filters.append(db.Worker.work_pool_id.in_(self.any_)) 

2093 return filters 

2094 

2095 

2096class WorkerFilterStatus(PrefectFilterBaseModel): 

2097 """Filter by `Worker.status`.""" 

2098 

2099 any_: Optional[list[schemas.statuses.WorkerStatus]] = Field( 

2100 default=None, description="A list of worker statuses to include" 

2101 ) 

2102 not_any_: Optional[list[schemas.statuses.WorkerStatus]] = Field( 

2103 default=None, description="A list of worker statuses to exclude" 

2104 ) 

2105 

2106 def _get_filter_list( 

2107 self, db: "PrefectDBInterface" 

2108 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

2109 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

2110 if self.any_ is not None: 2110 ↛ 2111line 2110 didn't jump to line 2111 because the condition on line 2110 was never true

2111 filters.append(db.Worker.status.in_(self.any_)) 

2112 if self.not_any_: 2112 ↛ 2113line 2112 didn't jump to line 2113 because the condition on line 2112 was never true

2113 filters.append(db.Worker.status.notin_(self.not_any_)) 

2114 return filters 

2115 

2116 

2117class WorkerFilterLastHeartbeatTime(PrefectFilterBaseModel): 

2118 """Filter by `Worker.last_heartbeat_time`.""" 

2119 

2120 before_: Optional[DateTime] = Field( 

2121 default=None, 

2122 description=( 

2123 "Only include processes whose last heartbeat was at or before this time" 

2124 ), 

2125 ) 

2126 after_: Optional[DateTime] = Field( 

2127 default=None, 

2128 description=( 

2129 "Only include processes whose last heartbeat was at or after this time" 

2130 ), 

2131 ) 

2132 

2133 def _get_filter_list( 

2134 self, db: "PrefectDBInterface" 

2135 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

2136 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

2137 if self.before_ is not None: 

2138 filters.append(db.Worker.last_heartbeat_time <= self.before_) 

2139 if self.after_ is not None: 

2140 filters.append(db.Worker.last_heartbeat_time >= self.after_) 

2141 return filters 

2142 

2143 

2144class WorkerFilter(PrefectOperatorFilterBaseModel): 

2145 """Filter by `Worker.last_heartbeat_time`.""" 

2146 

2147 # worker_config_id: Optional[WorkerFilterWorkPoolId] = Field( 

2148 # default=None, description="Filter criteria for `Worker.worker_config_id`" 

2149 # ) 

2150 

2151 last_heartbeat_time: Optional[WorkerFilterLastHeartbeatTime] = Field( 

2152 default=None, 

2153 description="Filter criteria for `Worker.last_heartbeat_time`", 

2154 ) 

2155 

2156 status: Optional[WorkerFilterStatus] = Field( 

2157 default=None, description="Filter criteria for `Worker.status`" 

2158 ) 

2159 

2160 def _get_filter_list( 

2161 self, db: "PrefectDBInterface" 

2162 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

2163 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

2164 

2165 if self.last_heartbeat_time is not None: 2165 ↛ 2166line 2165 didn't jump to line 2166 because the condition on line 2165 was never true

2166 filters.append(self.last_heartbeat_time.as_sql_filter()) 

2167 

2168 if self.status is not None: 

2169 filters.append(self.status.as_sql_filter()) 

2170 

2171 return filters 

2172 

2173 

2174class ArtifactFilterId(PrefectFilterBaseModel): 

2175 """Filter by `Artifact.id`.""" 

2176 

2177 any_: Optional[list[UUID]] = Field( 

2178 default=None, description="A list of artifact ids to include" 

2179 ) 

2180 

2181 def _get_filter_list( 

2182 self, db: "PrefectDBInterface" 

2183 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

2184 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

2185 if self.any_ is not None: 

2186 filters.append(db.Artifact.id.in_(self.any_)) 

2187 return filters 

2188 

2189 

2190class ArtifactFilterKey(PrefectFilterBaseModel): 

2191 """Filter by `Artifact.key`.""" 

2192 

2193 any_: Optional[list[str]] = Field( 

2194 default=None, description="A list of artifact keys to include" 

2195 ) 

2196 

2197 like_: Optional[str] = Field( 

2198 default=None, 

2199 description=( 

2200 "A string to match artifact keys against. This can include " 

2201 "SQL wildcard characters like `%` and `_`." 

2202 ), 

2203 examples=["my-artifact-%"], 

2204 ) 

2205 

2206 exists_: Optional[bool] = Field( 

2207 default=None, 

2208 description=( 

2209 "If `true`, only include artifacts with a non-null key. If `false`, " 

2210 "only include artifacts with a null key." 

2211 ), 

2212 ) 

2213 

2214 def _get_filter_list( 

2215 self, db: "PrefectDBInterface" 

2216 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

2217 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

2218 if self.any_ is not None: 

2219 filters.append(db.Artifact.key.in_(self.any_)) 

2220 if self.like_: 

2221 filters.append(db.Artifact.key.ilike(f"%{self.like_}%")) 

2222 if self.exists_ is not None: 

2223 filters.append( 

2224 db.Artifact.key.isnot(None) 

2225 if self.exists_ 

2226 else db.Artifact.key.is_(None) 

2227 ) 

2228 return filters 

2229 

2230 

2231class ArtifactFilterFlowRunId(PrefectFilterBaseModel): 

2232 """Filter by `Artifact.flow_run_id`.""" 

2233 

2234 any_: Optional[list[UUID]] = Field( 

2235 default=None, description="A list of flow run IDs to include" 

2236 ) 

2237 

2238 def _get_filter_list( 

2239 self, db: "PrefectDBInterface" 

2240 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

2241 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

2242 if self.any_ is not None: 

2243 filters.append(db.Artifact.flow_run_id.in_(self.any_)) 

2244 return filters 

2245 

2246 

2247class ArtifactFilterTaskRunId(PrefectFilterBaseModel): 

2248 """Filter by `Artifact.task_run_id`.""" 

2249 

2250 any_: Optional[list[UUID]] = Field( 

2251 default=None, description="A list of task run IDs to include" 

2252 ) 

2253 

2254 def _get_filter_list( 

2255 self, db: "PrefectDBInterface" 

2256 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

2257 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

2258 if self.any_ is not None: 

2259 filters.append(db.Artifact.task_run_id.in_(self.any_)) 

2260 return filters 

2261 

2262 

2263class ArtifactFilterType(PrefectFilterBaseModel): 

2264 """Filter by `Artifact.type`.""" 

2265 

2266 any_: Optional[list[str]] = Field( 

2267 default=None, description="A list of artifact types to include" 

2268 ) 

2269 not_any_: Optional[list[str]] = Field( 

2270 default=None, description="A list of artifact types to exclude" 

2271 ) 

2272 

2273 def _get_filter_list( 

2274 self, db: "PrefectDBInterface" 

2275 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

2276 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

2277 if self.any_ is not None: 

2278 filters.append(db.Artifact.type.in_(self.any_)) 

2279 if self.not_any_: 

2280 filters.append(db.Artifact.type.notin_(self.not_any_)) 

2281 return filters 

2282 

2283 

2284class ArtifactFilter(PrefectOperatorFilterBaseModel): 

2285 """Filter artifacts. Only artifacts matching all criteria will be returned""" 

2286 

2287 id: Optional[ArtifactFilterId] = Field( 

2288 default=None, description="Filter criteria for `Artifact.id`" 

2289 ) 

2290 key: Optional[ArtifactFilterKey] = Field( 

2291 default=None, description="Filter criteria for `Artifact.key`" 

2292 ) 

2293 flow_run_id: Optional[ArtifactFilterFlowRunId] = Field( 

2294 default=None, description="Filter criteria for `Artifact.flow_run_id`" 

2295 ) 

2296 task_run_id: Optional[ArtifactFilterTaskRunId] = Field( 

2297 default=None, description="Filter criteria for `Artifact.task_run_id`" 

2298 ) 

2299 type: Optional[ArtifactFilterType] = Field( 

2300 default=None, description="Filter criteria for `Artifact.type`" 

2301 ) 

2302 

2303 def _get_filter_list( 

2304 self, db: "PrefectDBInterface" 

2305 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

2306 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

2307 

2308 if self.id is not None: 

2309 filters.append(self.id.as_sql_filter()) 

2310 if self.key is not None: 

2311 filters.append(self.key.as_sql_filter()) 

2312 if self.flow_run_id is not None: 

2313 filters.append(self.flow_run_id.as_sql_filter()) 

2314 if self.task_run_id is not None: 

2315 filters.append(self.task_run_id.as_sql_filter()) 

2316 if self.type is not None: 

2317 filters.append(self.type.as_sql_filter()) 

2318 

2319 return filters 

2320 

2321 

2322class ArtifactCollectionFilterLatestId(PrefectFilterBaseModel): 

2323 """Filter by `ArtifactCollection.latest_id`.""" 

2324 

2325 any_: Optional[list[UUID]] = Field( 

2326 default=None, description="A list of artifact ids to include" 

2327 ) 

2328 

2329 def _get_filter_list( 

2330 self, db: "PrefectDBInterface" 

2331 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

2332 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

2333 if self.any_ is not None: 

2334 filters.append(db.ArtifactCollection.latest_id.in_(self.any_)) 

2335 return filters 

2336 

2337 

2338class ArtifactCollectionFilterKey(PrefectFilterBaseModel): 

2339 """Filter by `ArtifactCollection.key`.""" 

2340 

2341 any_: Optional[list[str]] = Field( 

2342 default=None, description="A list of artifact keys to include" 

2343 ) 

2344 

2345 like_: Optional[str] = Field( 

2346 default=None, 

2347 description=( 

2348 "A string to match artifact keys against. This can include " 

2349 "SQL wildcard characters like `%` and `_`." 

2350 ), 

2351 examples=["my-artifact-%"], 

2352 ) 

2353 

2354 exists_: Optional[bool] = Field( 

2355 default=None, 

2356 description=( 

2357 "If `true`, only include artifacts with a non-null key. If `false`, " 

2358 "only include artifacts with a null key. Should return all rows in " 

2359 "the ArtifactCollection table if specified." 

2360 ), 

2361 ) 

2362 

2363 def _get_filter_list( 

2364 self, db: "PrefectDBInterface" 

2365 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

2366 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

2367 if self.any_ is not None: 

2368 filters.append(db.ArtifactCollection.key.in_(self.any_)) 

2369 if self.like_: 

2370 filters.append(db.ArtifactCollection.key.ilike(f"%{self.like_}%")) 

2371 if self.exists_ is not None: 

2372 filters.append( 

2373 db.ArtifactCollection.key.isnot(None) 

2374 if self.exists_ 

2375 else db.ArtifactCollection.key.is_(None) 

2376 ) 

2377 return filters 

2378 

2379 

2380class ArtifactCollectionFilterFlowRunId(PrefectFilterBaseModel): 

2381 """Filter by `ArtifactCollection.flow_run_id`.""" 

2382 

2383 any_: Optional[list[UUID]] = Field( 

2384 default=None, description="A list of flow run IDs to include" 

2385 ) 

2386 

2387 def _get_filter_list( 

2388 self, db: "PrefectDBInterface" 

2389 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

2390 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

2391 if self.any_ is not None: 

2392 filters.append(db.ArtifactCollection.flow_run_id.in_(self.any_)) 

2393 return filters 

2394 

2395 

2396class ArtifactCollectionFilterTaskRunId(PrefectFilterBaseModel): 

2397 """Filter by `ArtifactCollection.task_run_id`.""" 

2398 

2399 any_: Optional[list[UUID]] = Field( 

2400 default=None, description="A list of task run IDs to include" 

2401 ) 

2402 

2403 def _get_filter_list( 

2404 self, db: "PrefectDBInterface" 

2405 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

2406 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

2407 if self.any_ is not None: 

2408 filters.append(db.ArtifactCollection.task_run_id.in_(self.any_)) 

2409 return filters 

2410 

2411 

2412class ArtifactCollectionFilterType(PrefectFilterBaseModel): 

2413 """Filter by `ArtifactCollection.type`.""" 

2414 

2415 any_: Optional[list[str]] = Field( 

2416 default=None, description="A list of artifact types to include" 

2417 ) 

2418 not_any_: Optional[list[str]] = Field( 

2419 default=None, description="A list of artifact types to exclude" 

2420 ) 

2421 

2422 def _get_filter_list( 

2423 self, db: "PrefectDBInterface" 

2424 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

2425 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

2426 if self.any_ is not None: 

2427 filters.append(db.ArtifactCollection.type.in_(self.any_)) 

2428 if self.not_any_: 

2429 filters.append(db.ArtifactCollection.type.notin_(self.not_any_)) 

2430 return filters 

2431 

2432 

2433class ArtifactCollectionFilter(PrefectOperatorFilterBaseModel): 

2434 """Filter artifact collections. Only artifact collections matching all criteria will be returned""" 

2435 

2436 latest_id: Optional[ArtifactCollectionFilterLatestId] = Field( 

2437 default=None, description="Filter criteria for `Artifact.id`" 

2438 ) 

2439 key: Optional[ArtifactCollectionFilterKey] = Field( 

2440 default=None, description="Filter criteria for `Artifact.key`" 

2441 ) 

2442 flow_run_id: Optional[ArtifactCollectionFilterFlowRunId] = Field( 

2443 default=None, description="Filter criteria for `Artifact.flow_run_id`" 

2444 ) 

2445 task_run_id: Optional[ArtifactCollectionFilterTaskRunId] = Field( 

2446 default=None, description="Filter criteria for `Artifact.task_run_id`" 

2447 ) 

2448 type: Optional[ArtifactCollectionFilterType] = Field( 

2449 default=None, description="Filter criteria for `Artifact.type`" 

2450 ) 

2451 

2452 def _get_filter_list( 

2453 self, db: "PrefectDBInterface" 

2454 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

2455 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

2456 

2457 if self.latest_id is not None: 

2458 filters.append(self.latest_id.as_sql_filter()) 

2459 if self.key is not None: 

2460 filters.append(self.key.as_sql_filter()) 

2461 if self.flow_run_id is not None: 

2462 filters.append(self.flow_run_id.as_sql_filter()) 

2463 if self.task_run_id is not None: 

2464 filters.append(self.task_run_id.as_sql_filter()) 

2465 if self.type is not None: 

2466 filters.append(self.type.as_sql_filter()) 

2467 

2468 return filters 

2469 

2470 

2471class VariableFilterId(PrefectFilterBaseModel): 

2472 """Filter by `Variable.id`.""" 

2473 

2474 any_: Optional[list[UUID]] = Field( 

2475 default=None, description="A list of variable ids to include" 

2476 ) 

2477 

2478 def _get_filter_list( 

2479 self, db: "PrefectDBInterface" 

2480 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

2481 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

2482 if self.any_ is not None: 

2483 filters.append(db.Variable.id.in_(self.any_)) 

2484 return filters 

2485 

2486 

2487class VariableFilterName(PrefectFilterBaseModel): 

2488 """Filter by `Variable.name`.""" 

2489 

2490 any_: Optional[list[str]] = Field( 

2491 default=None, description="A list of variables names to include" 

2492 ) 

2493 like_: Optional[str] = Field( 

2494 default=None, 

2495 description=( 

2496 "A string to match variable names against. This can include " 

2497 "SQL wildcard characters like `%` and `_`." 

2498 ), 

2499 examples=["my_variable_%"], 

2500 ) 

2501 

2502 def _get_filter_list( 

2503 self, db: "PrefectDBInterface" 

2504 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

2505 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

2506 if self.any_ is not None: 

2507 filters.append(db.Variable.name.in_(self.any_)) 

2508 if self.like_: 

2509 filters.append(db.Variable.name.ilike(f"%{self.like_}%")) 

2510 return filters 

2511 

2512 

2513class VariableFilterTags(PrefectOperatorFilterBaseModel): 

2514 """Filter by `Variable.tags`.""" 

2515 

2516 all_: Optional[list[str]] = Field( 

2517 default=None, 

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

2519 description=( 

2520 "A list of tags. Variables will be returned only if their tags are a" 

2521 " superset of the list" 

2522 ), 

2523 ) 

2524 is_null_: Optional[bool] = Field( 

2525 default=None, description="If true, only include Variables without tags" 

2526 ) 

2527 

2528 def _get_filter_list( 

2529 self, db: "PrefectDBInterface" 

2530 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

2531 filters: list[sa.ColumnElement[bool]] = [] 

2532 if self.all_: 

2533 filters.append(db.Variable.tags.has_all(_as_array(self.all_))) 

2534 if self.is_null_ is not None: 

2535 filters.append( 

2536 db.Variable.tags == [] if self.is_null_ else db.Variable.tags != [] 

2537 ) 

2538 return filters 

2539 

2540 

2541class VariableFilter(PrefectOperatorFilterBaseModel): 

2542 """Filter variables. Only variables matching all criteria will be returned""" 

2543 

2544 id: Optional[VariableFilterId] = Field( 

2545 default=None, description="Filter criteria for `Variable.id`" 

2546 ) 

2547 name: Optional[VariableFilterName] = Field( 

2548 default=None, description="Filter criteria for `Variable.name`" 

2549 ) 

2550 tags: Optional[VariableFilterTags] = Field( 

2551 default=None, description="Filter criteria for `Variable.tags`" 

2552 ) 

2553 

2554 def _get_filter_list( 

2555 self, db: "PrefectDBInterface" 

2556 ) -> Iterable[sa.ColumnExpressionArgument[bool]]: 

2557 filters: list[sa.ColumnExpressionArgument[bool]] = [] 

2558 

2559 if self.id is not None: 

2560 filters.append(self.id.as_sql_filter()) 

2561 if self.name is not None: 

2562 filters.append(self.name.as_sql_filter()) 

2563 if self.tags is not None: 

2564 filters.append(self.tags.as_sql_filter()) 

2565 return filters