Coverage for /usr/local/lib/python3.12/site-packages/prefect/server/orchestration/policies.py: 97%
26 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"""
2Policies are collections of orchestration rules and transforms.
4Prefect implements (most) orchestration with logic that governs a Prefect flow or task
5changing state. Policies organize of orchestration logic both to provide an ordering
6mechanism as well as provide observability into the orchestration process.
8While Prefect's orchestration rules can gracefully run independently of one another, ordering can still have an impact on
9the observed behavior of the system. For example, it makes no sense to secure a
10concurrency slot for a run if a cached state exists. Furthermore, policies, provide a
11mechanism to configure and observe exactly what logic will fire against a transition.
12"""
14from __future__ import annotations
16from abc import ABC, abstractmethod
17from typing import Generic, TypeVar, Union
19from prefect.server.database import orm_models
20from prefect.server.orchestration.rules import (
21 BaseOrchestrationRule,
22 BaseUniversalTransform,
23)
24from prefect.server.schemas import core, states
26T = TypeVar("T", bound=orm_models.Run)
27RP = TypeVar("RP", bound=Union[core.FlowRunPolicy, core.TaskRunPolicy])
30class BaseOrchestrationPolicy(ABC, Generic[T, RP]):
31 """
32 An abstract base class used to organize orchestration rules in priority order.
34 Different collections of orchestration rules might be used to govern various kinds
35 of transitions. For example, flow-run states and task-run states might require
36 different orchestration logic.
37 """
39 @staticmethod
40 @abstractmethod
41 def priority() -> list[
42 type[BaseUniversalTransform[T, RP] | BaseOrchestrationRule[T, RP]]
43 ]:
44 """
45 A list of orchestration rules in priority order.
46 """
48 return []
50 @classmethod
51 def compile_transition_rules(
52 cls,
53 from_state: states.StateType | None = None,
54 to_state: states.StateType | None = None,
55 ) -> list[type[BaseUniversalTransform[T, RP] | BaseOrchestrationRule[T, RP]]]:
56 """
57 Returns rules in policy that are valid for the specified state transition.
58 """
60 transition_rules: list[
61 type[BaseUniversalTransform[T, RP] | BaseOrchestrationRule[T, RP]]
62 ] = []
63 for rule in cls.priority():
64 if from_state in rule.FROM_STATES and to_state in rule.TO_STATES:
65 transition_rules.append(rule)
66 return transition_rules
69class TaskRunOrchestrationPolicy(
70 BaseOrchestrationPolicy[orm_models.TaskRun, core.TaskRunPolicy]
71):
72 pass
75class FlowRunOrchestrationPolicy(
76 BaseOrchestrationPolicy[orm_models.FlowRun, core.FlowRunPolicy]
77):
78 pass
81class GenericOrchestrationPolicy(
82 BaseOrchestrationPolicy[
83 orm_models.Run, Union[core.FlowRunPolicy, core.TaskRunPolicy]
84 ]
85):
86 pass