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

1""" 

2Perpetual services are background services that run on a periodic schedule using docket. 

3 

4This module provides the registry and scheduling logic for perpetual services, 

5using docket's Perpetual dependency for distributed, HA-aware task scheduling. 

6""" 

7 

8from __future__ import annotations 

9 

10import logging 

11from dataclasses import dataclass 

12from typing import Callable, TypeVar 

13 

14from docket import Docket, Perpetual 

15from docket.dependencies import get_single_dependency_parameter_of_type 

16from docket.execution import TaskFunction 

17 

18from prefect.logging import get_logger 

19 

20logger: logging.Logger = get_logger(__name__) 

21 

22EnabledGetter = Callable[[], bool] 

23"""A callable that returns whether a service is enabled.""" 

24 

25F = TypeVar("F", bound=TaskFunction) 

26 

27 

28@dataclass 

29class PerpetualServiceConfig: 

30 """Configuration for a perpetual service function.""" 

31 

32 function: TaskFunction 

33 enabled_getter: EnabledGetter 

34 run_in_ephemeral: bool = False 

35 run_in_webserver: bool = False 

36 

37 

38# Registry of all perpetual service functions 

39_PERPETUAL_SERVICES: list[PerpetualServiceConfig] = [] 

40 

41 

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. 

49 

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. 

54 

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 """ 

62 

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 

73 

74 return decorator 

75 

76 

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. 

83 

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. 

87 

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 

99 

100 services.append(config) 

101 

102 return services 

103 

104 

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. 

111 

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. 

115 

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 ) 

127 

128 return services 

129 

130 

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. 

138 

139 Disabled services are not registered at all, so they never run. 

140 

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) 

148 

149 for config in enabled_services: 

150 docket.register(config.function) 

151 logger.debug(f"Registered perpetual service: {config.function.__name__}") 

152 

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 

161 

162 logger.info(f"Scheduling perpetual service: {config.function.__name__}") 

163 await docket.add(config.function, key=config.function.__name__)() 

164 

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 )