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
« 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.
4Each filter schema includes logic for transforming itself into a SQL `where` clause.
5"""
7from collections.abc import Iterable, Sequence
8from typing import TYPE_CHECKING, ClassVar, Optional
9from uuid import UUID
11from pydantic import ConfigDict, Field
12from sqlalchemy.dialects import postgresql
13from sqlalchemy.sql.functions import coalesce
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
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
27 from prefect.server.database import PrefectDBInterface
28 from prefect.server.schemas.core import Log
29else:
30 sa = lazy_import("sqlalchemy")
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
37def _as_array(elems: Sequence[str]) -> sa.ColumnElement[Sequence[str]]:
38 return sa.cast(postgresql.array(elems), type_=postgresql.ARRAY(sa.String()))
41class Operator(AutoEnum):
42 """Operators for combining filter criteria."""
44 and_ = AutoEnum.auto()
45 or_ = AutoEnum.auto()
48class PrefectFilterBaseModel(PrefectBaseModel):
49 """Base model for Prefect filters"""
51 model_config: ClassVar[ConfigDict] = ConfigDict(extra="forbid")
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
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)
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")
70class PrefectOperatorFilterBaseModel(PrefectFilterBaseModel):
71 """Base model for Prefect filters that combines criteria with a user-provided operator"""
73 operator: Operator = Field(
74 default=Operator.and_,
75 description="Operator for combining filter criteria. Defaults to 'and_'.",
76 )
78 def as_sql_filter(self) -> sa.ColumnElement[bool]:
79 from prefect.server.database.dependencies import provide_database_interface
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)
88class FlowFilterId(PrefectFilterBaseModel):
89 """Filter by `Flow.id`."""
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 )
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
109class FlowFilterDeployment(PrefectOperatorFilterBaseModel):
110 """Filter by flows by deployment"""
112 is_null_: Optional[bool] = Field(
113 default=None,
114 description="If true, only include flows without deployments",
115 )
117 def _get_filter_list(
118 self, db: "PrefectDBInterface"
119 ) -> Iterable[sa.ColumnExpressionArgument[bool]]:
120 filters: list[sa.ColumnExpressionArgument[bool]] = []
122 if self.is_null_ is not None:
123 deployments_subquery = (
124 sa.select(db.Deployment.flow_id).distinct().subquery()
125 )
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 )
136 return filters
139class FlowFilterName(PrefectFilterBaseModel):
140 """Filter by `Flow.name`."""
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 )
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 )
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
169class FlowFilterTags(PrefectOperatorFilterBaseModel):
170 """Filter by `Flow.tags`."""
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 )
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
195class FlowFilter(PrefectOperatorFilterBaseModel):
196 """Filter for flows. Only flows matching all criteria will be returned."""
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 )
211 def _get_filter_list(
212 self, db: "PrefectDBInterface"
213 ) -> Iterable[sa.ColumnExpressionArgument[bool]]:
214 filters: list[sa.ColumnExpressionArgument[bool]] = []
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())
225 return filters
228class FlowRunFilterId(PrefectFilterBaseModel):
229 """Filter by `FlowRun.id`."""
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 )
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
249class FlowRunFilterName(PrefectFilterBaseModel):
250 """Filter by `FlowRun.name`."""
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 )
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 )
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
279class FlowRunFilterTags(PrefectOperatorFilterBaseModel):
280 """Filter by `FlowRun.tags`."""
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 )
291 any_: Optional[list[str]] = Field(
292 default=None,
293 examples=[["tag-1", "tag-2"]],
294 description="A list of tags to include",
295 )
297 is_null_: Optional[bool] = Field(
298 default=None, description="If true, only include flow runs without tags"
299 )
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()))
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
319class FlowRunFilterDeploymentId(PrefectOperatorFilterBaseModel):
320 """Filter by `FlowRun.deployment_id`."""
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 )
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
345class FlowRunFilterWorkQueueName(PrefectOperatorFilterBaseModel):
346 """Filter by `FlowRun.work_queue_name`."""
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 )
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
373class FlowRunFilterStateType(PrefectFilterBaseModel):
374 """Filter by `FlowRun.state_type`."""
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 )
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
394class FlowRunFilterStateName(PrefectFilterBaseModel):
395 """Filter by `FlowRun.state_name`."""
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 )
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
415class FlowRunFilterState(PrefectOperatorFilterBaseModel):
416 """Filter by `FlowRun.state_type` and `FlowRun.state_name`."""
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 )
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
440class FlowRunFilterFlowVersion(PrefectFilterBaseModel):
441 """Filter by `FlowRun.flow_version`."""
443 any_: Optional[list[str]] = Field(
444 default=None, description="A list of flow run flow_versions to include"
445 )
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
456class FlowRunFilterStartTime(PrefectFilterBaseModel):
457 """Filter by `FlowRun.start_time`."""
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 )
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
494class FlowRunFilterEndTime(PrefectFilterBaseModel):
495 """Filter by `FlowRun.end_time`."""
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 )
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
526class FlowRunFilterExpectedStartTime(PrefectFilterBaseModel):
527 """Filter by `FlowRun.expected_start_time`."""
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 )
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
549class FlowRunFilterNextScheduledStartTime(PrefectFilterBaseModel):
550 """Filter by `FlowRun.next_scheduled_start_time`."""
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 )
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
578class FlowRunFilterParentFlowRunId(PrefectOperatorFilterBaseModel):
579 """Filter for subflows of a given flow run"""
581 any_: Optional[list[UUID]] = Field(
582 default=None, description="A list of parent flow run ids to include"
583 )
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
605class FlowRunFilterParentTaskRunId(PrefectOperatorFilterBaseModel):
606 """Filter by `FlowRun.parent_task_run_id`."""
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 )
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
631class FlowRunFilterIdempotencyKey(PrefectFilterBaseModel):
632 """Filter by FlowRun.idempotency_key."""
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 )
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
652class FlowRunFilterCreatedBy(PrefectOperatorFilterBaseModel):
653 """Filter by `FlowRun.created_by`."""
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 )
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
691class FlowRunFilter(PrefectOperatorFilterBaseModel):
692 """Filter flow runs. Only flow runs matching all criteria will be returned"""
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 )
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 )
761 def _get_filter_list(
762 self, db: "PrefectDBInterface"
763 ) -> Iterable[sa.ColumnExpressionArgument[bool]]:
764 filters: list[sa.ColumnExpressionArgument[bool]] = []
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())
797 return filters
800class TaskRunFilterFlowRunId(PrefectOperatorFilterBaseModel):
801 """Filter by `TaskRun.flow_run_id`."""
803 any_: Optional[list[UUID]] = Field(
804 default=None, description="A list of task run flow run ids to include"
805 )
807 is_null_: Optional[bool] = Field(
808 default=False, description="Filter for task runs with None as their flow run id"
809 )
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
825class TaskRunFilterId(PrefectFilterBaseModel):
826 """Filter by `TaskRun.id`."""
828 any_: Optional[list[UUID]] = Field(
829 default=None, description="A list of task run ids to include"
830 )
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
841class TaskRunFilterName(PrefectFilterBaseModel):
842 """Filter by `TaskRun.name`."""
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 )
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 )
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
871class TaskRunFilterTags(PrefectOperatorFilterBaseModel):
872 """Filter by `TaskRun.tags`."""
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 )
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
899class TaskRunFilterStateType(PrefectFilterBaseModel):
900 """Filter by `TaskRun.state_type`."""
902 any_: Optional[list[schemas.states.StateType]] = Field(
903 default=None, description="A list of task run state types to include"
904 )
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
915class TaskRunFilterStateName(PrefectFilterBaseModel):
916 """Filter by `TaskRun.state_name`."""
918 any_: Optional[list[str]] = Field(
919 default=None, description="A list of task run state names to include"
920 )
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
931class TaskRunFilterState(PrefectOperatorFilterBaseModel):
932 """Filter by `TaskRun.type` and `TaskRun.name`."""
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 )
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
956class TaskRunFilterSubFlowRuns(PrefectFilterBaseModel):
957 """Filter by `TaskRun.subflow_run`."""
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 )
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
978class TaskRunFilterStartTime(PrefectFilterBaseModel):
979 """Filter by `TaskRun.start_time`."""
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 )
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
1010class TaskRunFilterEndTime(PrefectFilterBaseModel):
1011 """Filter by `TaskRun.end_time`."""
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 )
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
1042class TaskRunFilterExpectedStartTime(PrefectFilterBaseModel):
1043 """Filter by `TaskRun.expected_start_time`."""
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 )
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
1065class TaskRunFilter(PrefectOperatorFilterBaseModel):
1066 """Filter task runs. Only task runs matching all criteria will be returned"""
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 )
1096 def _get_filter_list(
1097 self, db: "PrefectDBInterface"
1098 ) -> Iterable[sa.ColumnExpressionArgument[bool]]:
1099 filters: list[sa.ColumnExpressionArgument[bool]] = []
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())
1120 return filters
1123class DeploymentFilterId(PrefectFilterBaseModel):
1124 """Filter by `Deployment.id`."""
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 )
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
1144class DeploymentFilterName(PrefectFilterBaseModel):
1145 """Filter by `Deployment.name`."""
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 )
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 )
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
1174class DeploymentOrFlowNameFilter(PrefectFilterBaseModel):
1175 """Filter by `Deployment.name` or `Flow.name` with a single input string for ilike filtering."""
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 )
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_}%")
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
1199class DeploymentFilterPaused(PrefectFilterBaseModel):
1200 """Filter by `Deployment.paused`."""
1202 eq_: Optional[bool] = Field(
1203 default=None,
1204 description="Only returns where deployment is/is not paused",
1205 )
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
1216class DeploymentFilterWorkQueueName(PrefectFilterBaseModel):
1217 """Filter by `Deployment.work_queue_name`."""
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 )
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
1234class DeploymentFilterConcurrencyLimit(PrefectFilterBaseModel):
1235 """DEPRECATED: Prefer `Deployment.concurrency_limit_id` over `Deployment.concurrency_limit`."""
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 )
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 )
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 []
1259class DeploymentFilterTags(PrefectOperatorFilterBaseModel):
1260 """Filter by `Deployment.tags`."""
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 )
1276 is_null_: Optional[bool] = Field(
1277 default=None, description="If true, only include deployments without tags"
1278 )
1280 def _get_filter_list(
1281 self, db: "PrefectDBInterface"
1282 ) -> Iterable[sa.ColumnExpressionArgument[bool]]:
1283 from prefect.server.database import orm_models
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
1297class DeploymentFilter(PrefectOperatorFilterBaseModel):
1298 """Filter for deployments. Only deployments matching all criteria will be returned."""
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 )
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())
1341 return filters
1344class DeploymentScheduleFilterActive(PrefectFilterBaseModel):
1345 """Filter by `DeploymentSchedule.active`."""
1347 eq_: Optional[bool] = Field(
1348 default=None,
1349 description="Only returns where deployment schedule is/is not active",
1350 )
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
1361class DeploymentScheduleFilter(PrefectOperatorFilterBaseModel):
1362 """Filter for deployments. Only deployments matching all criteria will be returned."""
1364 active: Optional[DeploymentScheduleFilterActive] = Field(
1365 default=None, description="Filter criteria for `DeploymentSchedule.active`"
1366 )
1368 def _get_filter_list(
1369 self, db: "PrefectDBInterface"
1370 ) -> Iterable[sa.ColumnExpressionArgument[bool]]:
1371 filters: list[sa.ColumnExpressionArgument[bool]] = []
1373 if self.active is not None:
1374 filters.append(self.active.as_sql_filter())
1376 return filters
1379class LogFilterName(PrefectFilterBaseModel):
1380 """Filter by `Log.name`."""
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 )
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
1397class LogFilterLevel(PrefectFilterBaseModel):
1398 """Filter by `Log.level`."""
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 )
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 )
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
1423class LogFilterTimestamp(PrefectFilterBaseModel):
1424 """Filter by `Log.timestamp`."""
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 )
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
1446class LogFilterFlowRunId(PrefectFilterBaseModel):
1447 """Filter by `Log.flow_run_id`."""
1449 any_: Optional[list[UUID]] = Field(
1450 default=None, description="A list of flow run IDs to include"
1451 )
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
1462class LogFilterTaskRunId(PrefectFilterBaseModel):
1463 """Filter by `Log.task_run_id`."""
1465 any_: Optional[list[UUID]] = Field(
1466 default=None, description="A list of task run IDs to include"
1467 )
1469 is_null_: Optional[bool] = Field(
1470 default=None,
1471 description="If true, only include logs without a task run id",
1472 )
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
1489class LogFilterTextSearch(PrefectFilterBaseModel):
1490 """Filter by text search across log content."""
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 )
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
1507 if not isinstance(log, Log):
1508 raise TypeError(f"Expected Log object, got {type(log)}")
1510 # Parse query into components
1511 parsed = parse_text_search_query(self.query)
1513 # Build searchable text from message and logger name
1514 searchable_text = f"{log.message} {log.name}".lower()
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
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
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
1540 return True
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]] = []
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
1551 parsed = parse_text_search_query(self.query)
1553 # Build combined searchable text field (message + name)
1554 searchable_field = sa.func.concat(db.Log.message, " ", db.Log.name)
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 )
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))
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 )
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))
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 )
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))
1589 return filters
1592class LogFilter(PrefectOperatorFilterBaseModel):
1593 """Filter logs. Only logs matching all criteria will be returned"""
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 )
1611 def _get_filter_list(
1612 self, db: "PrefectDBInterface"
1613 ) -> Iterable[sa.ColumnExpressionArgument[bool]]:
1614 filters: list[sa.ColumnExpressionArgument[bool]] = []
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))
1627 return filters
1630class FilterSet(PrefectBaseModel):
1631 """A collection of filters for common objects"""
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 )
1648class BlockTypeFilterName(PrefectFilterBaseModel):
1649 """Filter by `BlockType.name`"""
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 )
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
1670class BlockTypeFilterSlug(PrefectFilterBaseModel):
1671 """Filter by `BlockType.slug`"""
1673 any_: Optional[list[str]] = Field(
1674 default=None, description="A list of slugs to match"
1675 )
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_))
1684 return filters
1687class BlockTypeFilter(PrefectFilterBaseModel):
1688 """Filter BlockTypes"""
1690 name: Optional[BlockTypeFilterName] = Field(
1691 default=None, description="Filter criteria for `BlockType.name`"
1692 )
1694 slug: Optional[BlockTypeFilterSlug] = Field(
1695 default=None, description="Filter criteria for `BlockType.slug`"
1696 )
1698 def _get_filter_list(
1699 self, db: "PrefectDBInterface"
1700 ) -> Iterable[sa.ColumnExpressionArgument[bool]]:
1701 filters: list[sa.ColumnExpressionArgument[bool]] = []
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())
1708 return filters
1711class BlockSchemaFilterBlockTypeId(PrefectFilterBaseModel):
1712 """Filter by `BlockSchema.block_type_id`."""
1714 any_: Optional[list[UUID]] = Field(
1715 default=None, description="A list of block type ids to include"
1716 )
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
1727class BlockSchemaFilterId(PrefectFilterBaseModel):
1728 """Filter by BlockSchema.id"""
1730 any_: Optional[list[UUID]] = Field(
1731 default=None, description="A list of IDs to include"
1732 )
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
1743class BlockSchemaFilterCapabilities(PrefectFilterBaseModel):
1744 """Filter by `BlockSchema.capabilities`"""
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 )
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
1764class BlockSchemaFilterVersion(PrefectFilterBaseModel):
1765 """Filter by `BlockSchema.capabilities`"""
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 )
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
1782class BlockSchemaFilter(PrefectOperatorFilterBaseModel):
1783 """Filter BlockSchemas"""
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 )
1798 def _get_filter_list(
1799 self, db: "PrefectDBInterface"
1800 ) -> Iterable[sa.ColumnExpressionArgument[bool]]:
1801 filters: list[sa.ColumnExpressionArgument[bool]] = []
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())
1812 return filters
1815class BlockDocumentFilterIsAnonymous(PrefectFilterBaseModel):
1816 """Filter by `BlockDocument.is_anonymous`."""
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 )
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
1834class BlockDocumentFilterBlockTypeId(PrefectFilterBaseModel):
1835 """Filter by `BlockDocument.block_type_id`."""
1837 any_: Optional[list[UUID]] = Field(
1838 default=None, description="A list of block type ids to include"
1839 )
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
1850class BlockDocumentFilterId(PrefectFilterBaseModel):
1851 """Filter by `BlockDocument.id`."""
1853 any_: Optional[list[UUID]] = Field(
1854 default=None, description="A list of block ids to include"
1855 )
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
1866class BlockDocumentFilterName(PrefectFilterBaseModel):
1867 """Filter by `BlockDocument.name`."""
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 )
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
1892class BlockDocumentFilter(PrefectOperatorFilterBaseModel):
1893 """Filter BlockDocuments. Only BlockDocuments matching all criteria will be returned"""
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 )
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
1928class WorkQueueFilterId(PrefectFilterBaseModel):
1929 """Filter by `WorkQueue.id`."""
1931 any_: Optional[list[UUID]] = Field(
1932 default=None,
1933 description="A list of work queue ids to include",
1934 )
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
1945class WorkQueueFilterName(PrefectFilterBaseModel):
1946 """Filter by `WorkQueue.name`."""
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 )
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 )
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
1979class WorkQueueFilter(PrefectOperatorFilterBaseModel):
1980 """Filter work queues. Only work queues matching all criteria will be
1981 returned"""
1983 id: Optional[WorkQueueFilterId] = Field(
1984 default=None, description="Filter criteria for `WorkQueue.id`"
1985 )
1987 name: Optional[WorkQueueFilterName] = Field(
1988 default=None, description="Filter criteria for `WorkQueue.name`"
1989 )
1991 def _get_filter_list(
1992 self, db: "PrefectDBInterface"
1993 ) -> Iterable[sa.ColumnExpressionArgument[bool]]:
1994 filters: list[sa.ColumnExpressionArgument[bool]] = []
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())
2001 return filters
2004class WorkPoolFilterId(PrefectFilterBaseModel):
2005 """Filter by `WorkPool.id`."""
2007 any_: Optional[list[UUID]] = Field(
2008 default=None, description="A list of work pool ids to include"
2009 )
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
2020class WorkPoolFilterName(PrefectFilterBaseModel):
2021 """Filter by `WorkPool.name`."""
2023 any_: Optional[list[str]] = Field(
2024 default=None, description="A list of work pool names to include"
2025 )
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
2036class WorkPoolFilterType(PrefectFilterBaseModel):
2037 """Filter by `WorkPool.type`."""
2039 any_: Optional[list[str]] = Field(
2040 default=None, description="A list of work pool types to include"
2041 )
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
2052class WorkPoolFilter(PrefectOperatorFilterBaseModel):
2053 """Filter work pools. Only work pools matching all criteria will be returned"""
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 )
2065 def _get_filter_list(
2066 self, db: "PrefectDBInterface"
2067 ) -> Iterable[sa.ColumnExpressionArgument[bool]]:
2068 filters: list[sa.ColumnExpressionArgument[bool]] = []
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())
2077 return filters
2080class WorkerFilterWorkPoolId(PrefectFilterBaseModel):
2081 """Filter by `Worker.worker_config_id`."""
2083 any_: Optional[list[UUID]] = Field(
2084 default=None, description="A list of work pool ids to include"
2085 )
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
2096class WorkerFilterStatus(PrefectFilterBaseModel):
2097 """Filter by `Worker.status`."""
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 )
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
2117class WorkerFilterLastHeartbeatTime(PrefectFilterBaseModel):
2118 """Filter by `Worker.last_heartbeat_time`."""
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 )
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
2144class WorkerFilter(PrefectOperatorFilterBaseModel):
2145 """Filter by `Worker.last_heartbeat_time`."""
2147 # worker_config_id: Optional[WorkerFilterWorkPoolId] = Field(
2148 # default=None, description="Filter criteria for `Worker.worker_config_id`"
2149 # )
2151 last_heartbeat_time: Optional[WorkerFilterLastHeartbeatTime] = Field(
2152 default=None,
2153 description="Filter criteria for `Worker.last_heartbeat_time`",
2154 )
2156 status: Optional[WorkerFilterStatus] = Field(
2157 default=None, description="Filter criteria for `Worker.status`"
2158 )
2160 def _get_filter_list(
2161 self, db: "PrefectDBInterface"
2162 ) -> Iterable[sa.ColumnExpressionArgument[bool]]:
2163 filters: list[sa.ColumnExpressionArgument[bool]] = []
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())
2168 if self.status is not None:
2169 filters.append(self.status.as_sql_filter())
2171 return filters
2174class ArtifactFilterId(PrefectFilterBaseModel):
2175 """Filter by `Artifact.id`."""
2177 any_: Optional[list[UUID]] = Field(
2178 default=None, description="A list of artifact ids to include"
2179 )
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
2190class ArtifactFilterKey(PrefectFilterBaseModel):
2191 """Filter by `Artifact.key`."""
2193 any_: Optional[list[str]] = Field(
2194 default=None, description="A list of artifact keys to include"
2195 )
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 )
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 )
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
2231class ArtifactFilterFlowRunId(PrefectFilterBaseModel):
2232 """Filter by `Artifact.flow_run_id`."""
2234 any_: Optional[list[UUID]] = Field(
2235 default=None, description="A list of flow run IDs to include"
2236 )
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
2247class ArtifactFilterTaskRunId(PrefectFilterBaseModel):
2248 """Filter by `Artifact.task_run_id`."""
2250 any_: Optional[list[UUID]] = Field(
2251 default=None, description="A list of task run IDs to include"
2252 )
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
2263class ArtifactFilterType(PrefectFilterBaseModel):
2264 """Filter by `Artifact.type`."""
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 )
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
2284class ArtifactFilter(PrefectOperatorFilterBaseModel):
2285 """Filter artifacts. Only artifacts matching all criteria will be returned"""
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 )
2303 def _get_filter_list(
2304 self, db: "PrefectDBInterface"
2305 ) -> Iterable[sa.ColumnExpressionArgument[bool]]:
2306 filters: list[sa.ColumnExpressionArgument[bool]] = []
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())
2319 return filters
2322class ArtifactCollectionFilterLatestId(PrefectFilterBaseModel):
2323 """Filter by `ArtifactCollection.latest_id`."""
2325 any_: Optional[list[UUID]] = Field(
2326 default=None, description="A list of artifact ids to include"
2327 )
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
2338class ArtifactCollectionFilterKey(PrefectFilterBaseModel):
2339 """Filter by `ArtifactCollection.key`."""
2341 any_: Optional[list[str]] = Field(
2342 default=None, description="A list of artifact keys to include"
2343 )
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 )
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 )
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
2380class ArtifactCollectionFilterFlowRunId(PrefectFilterBaseModel):
2381 """Filter by `ArtifactCollection.flow_run_id`."""
2383 any_: Optional[list[UUID]] = Field(
2384 default=None, description="A list of flow run IDs to include"
2385 )
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
2396class ArtifactCollectionFilterTaskRunId(PrefectFilterBaseModel):
2397 """Filter by `ArtifactCollection.task_run_id`."""
2399 any_: Optional[list[UUID]] = Field(
2400 default=None, description="A list of task run IDs to include"
2401 )
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
2412class ArtifactCollectionFilterType(PrefectFilterBaseModel):
2413 """Filter by `ArtifactCollection.type`."""
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 )
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
2433class ArtifactCollectionFilter(PrefectOperatorFilterBaseModel):
2434 """Filter artifact collections. Only artifact collections matching all criteria will be returned"""
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 )
2452 def _get_filter_list(
2453 self, db: "PrefectDBInterface"
2454 ) -> Iterable[sa.ColumnExpressionArgument[bool]]:
2455 filters: list[sa.ColumnExpressionArgument[bool]] = []
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())
2468 return filters
2471class VariableFilterId(PrefectFilterBaseModel):
2472 """Filter by `Variable.id`."""
2474 any_: Optional[list[UUID]] = Field(
2475 default=None, description="A list of variable ids to include"
2476 )
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
2487class VariableFilterName(PrefectFilterBaseModel):
2488 """Filter by `Variable.name`."""
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 )
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
2513class VariableFilterTags(PrefectOperatorFilterBaseModel):
2514 """Filter by `Variable.tags`."""
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 )
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
2541class VariableFilter(PrefectOperatorFilterBaseModel):
2542 """Filter variables. Only variables matching all criteria will be returned"""
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 )
2554 def _get_filter_list(
2555 self, db: "PrefectDBInterface"
2556 ) -> Iterable[sa.ColumnExpressionArgument[bool]]:
2557 filters: list[sa.ColumnExpressionArgument[bool]] = []
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