Coverage for /usr/local/lib/python3.12/site-packages/prefect/server/services/perpetual_services.py: 84%
59 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"""
2Perpetual services are background services that run on a periodic schedule using docket.
4This module provides the registry and scheduling logic for perpetual services,
5using docket's Perpetual dependency for distributed, HA-aware task scheduling.
6"""
8from __future__ import annotations
10import logging
11from dataclasses import dataclass
12from typing import Callable, TypeVar
14from docket import Docket, Perpetual
15from docket.dependencies import get_single_dependency_parameter_of_type
16from docket.execution import TaskFunction
18from prefect.logging import get_logger
20logger: logging.Logger = get_logger(__name__)
22EnabledGetter = Callable[[], bool]
23"""A callable that returns whether a service is enabled."""
25F = TypeVar("F", bound=TaskFunction)
28@dataclass
29class PerpetualServiceConfig:
30 """Configuration for a perpetual service function."""
32 function: TaskFunction
33 enabled_getter: EnabledGetter
34 run_in_ephemeral: bool = False
35 run_in_webserver: bool = False
38# Registry of all perpetual service functions
39_PERPETUAL_SERVICES: list[PerpetualServiceConfig] = []
42def perpetual_service(
43 enabled_getter: EnabledGetter,
44 run_in_ephemeral: bool = False,
45 run_in_webserver: bool = False,
46) -> Callable[[F], F]:
47 """
48 Decorator to register a perpetual service function.
50 Args:
51 enabled_getter: A callable that returns whether the service is enabled.
52 run_in_ephemeral: If True, this service runs in ephemeral server mode.
53 run_in_webserver: If True, this service runs in webserver-only mode.
55 Example:
56 @perpetual_service(
57 enabled_getter=lambda: get_current_settings().server.services.scheduler.enabled,
58 )
59 async def schedule_deployments(...) -> None:
60 ...
61 """
63 def decorator(func: F) -> F:
64 _PERPETUAL_SERVICES.append(
65 PerpetualServiceConfig(
66 function=func,
67 enabled_getter=enabled_getter,
68 run_in_ephemeral=run_in_ephemeral,
69 run_in_webserver=run_in_webserver,
70 )
71 )
72 return func
74 return decorator
77def get_perpetual_services(
78 ephemeral: bool = False,
79 webserver_only: bool = False,
80) -> list[PerpetualServiceConfig]:
81 """
82 Get perpetual services that should run in the current mode.
84 Args:
85 ephemeral: If True, only return services marked with run_in_ephemeral.
86 webserver_only: If True, only return services marked with run_in_webserver.
88 Returns:
89 List of perpetual service configurations to run.
90 """
91 services = []
92 for config in _PERPETUAL_SERVICES:
93 if webserver_only: 93 ↛ 94line 93 didn't jump to line 94 because the condition on line 93 was never true
94 if not config.run_in_webserver:
95 continue
96 elif ephemeral: 96 ↛ 97line 96 didn't jump to line 97 because the condition on line 96 was never true
97 if not config.run_in_ephemeral:
98 continue
100 services.append(config)
102 return services
105def get_enabled_perpetual_services(
106 ephemeral: bool = False,
107 webserver_only: bool = False,
108) -> list[PerpetualServiceConfig]:
109 """
110 Get perpetual services that are enabled and should run in the current mode.
112 Args:
113 ephemeral: If True, only return services marked with run_in_ephemeral.
114 webserver_only: If True, only return services marked with run_in_webserver.
116 Returns:
117 List of enabled perpetual service configurations.
118 """
119 services = []
120 for config in get_perpetual_services(ephemeral, webserver_only):
121 if config.enabled_getter():
122 services.append(config)
123 else:
124 logger.debug(
125 f"Skipping disabled perpetual service: {config.function.__name__}"
126 )
128 return services
131async def register_and_schedule_perpetual_services(
132 docket: Docket,
133 ephemeral: bool = False,
134 webserver_only: bool = False,
135) -> None:
136 """
137 Register enabled perpetual service functions with docket and schedule them.
139 Disabled services are not registered at all, so they never run.
141 Args:
142 docket: The docket instance to register functions with.
143 ephemeral: If True, only register services for ephemeral mode.
144 webserver_only: If True, only register services for webserver mode.
145 """
146 all_services = get_perpetual_services(ephemeral, webserver_only)
147 enabled_services = get_enabled_perpetual_services(ephemeral, webserver_only)
149 for config in enabled_services:
150 docket.register(config.function)
151 logger.debug(f"Registered perpetual service: {config.function.__name__}")
153 for config in enabled_services:
154 perpetual = get_single_dependency_parameter_of_type(config.function, Perpetual)
155 if perpetual is None: 155 ↛ 156line 155 didn't jump to line 156 because the condition on line 155 was never true
156 logger.warning(
157 f"Perpetual service {config.function.__name__} has no Perpetual "
158 "dependency - skipping scheduling"
159 )
160 continue
162 logger.info(f"Scheduling perpetual service: {config.function.__name__}")
163 await docket.add(config.function, key=config.function.__name__)()
165 total = len(all_services)
166 enabled = len(enabled_services)
167 disabled = total - enabled
168 logger.info(
169 f"Perpetual services: {enabled} enabled, {disabled} disabled, {total} total"
170 )