Coverage for /usr/local/lib/python3.12/site-packages/prefect/server/services/foreman.py: 96%

53 statements  

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

1""" 

2The Foreman service. Monitors workers and marks stale resources as offline/not ready. 

3""" 

4 

5from __future__ import annotations 

6 

7import logging 

8from datetime import timedelta 

9 

10import sqlalchemy as sa 

11from docket import Perpetual 

12 

13from prefect.logging import get_logger 

14from prefect.server import models 

15from prefect.server.database import PrefectDBInterface, provide_database_interface 

16from prefect.server.models.deployments import mark_deployments_not_ready 

17from prefect.server.models.work_queues import mark_work_queues_not_ready 

18from prefect.server.models.workers import emit_work_pool_status_event 

19from prefect.server.schemas.internal import InternalWorkPoolUpdate 

20from prefect.server.schemas.statuses import ( 

21 DeploymentStatus, 

22 WorkerStatus, 

23 WorkPoolStatus, 

24) 

25from prefect.server.services.perpetual_services import perpetual_service 

26from prefect.settings.context import get_current_settings 

27from prefect.types._datetime import now 

28 

29logger: logging.Logger = get_logger(__name__) 

30 

31 

32@perpetual_service( 

33 enabled_getter=lambda: get_current_settings().server.services.foreman.enabled, 

34) 

35async def monitor_worker_health( 

36 perpetual: Perpetual = Perpetual( 

37 automatic=True, 

38 every=timedelta( 

39 seconds=get_current_settings().server.services.foreman.loop_seconds 

40 ), 

41 ), 

42) -> None: 

43 """ 

44 Monitor workers and mark stale resources as offline/not ready. 

45 

46 Iterates over workers currently marked as online. Marks workers as offline 

47 if they have an old last_heartbeat_time. Marks work pools as not ready 

48 if they do not have any online workers and are currently marked as ready. 

49 Marks deployments as not ready if they have a last_polled time that is 

50 older than the configured deployment last polled timeout. 

51 """ 

52 settings = get_current_settings().server.services.foreman 

53 db = provide_database_interface() 

54 

55 await _mark_online_workers_without_recent_heartbeat_as_offline( 

56 db=db, 

57 inactivity_heartbeat_multiple=settings.inactivity_heartbeat_multiple, 

58 fallback_heartbeat_interval_seconds=settings.fallback_heartbeat_interval_seconds, 

59 ) 

60 await _mark_work_pools_as_not_ready(db=db) 

61 await _mark_deployments_as_not_ready( 

62 db=db, 

63 deployment_last_polled_timeout_seconds=settings.deployment_last_polled_timeout_seconds, 

64 ) 

65 await _mark_work_queues_as_not_ready( 

66 db=db, 

67 work_queue_last_polled_timeout_seconds=settings.work_queue_last_polled_timeout_seconds, 

68 ) 

69 

70 

71async def _mark_online_workers_without_recent_heartbeat_as_offline( 

72 db: PrefectDBInterface, 

73 inactivity_heartbeat_multiple: int, 

74 fallback_heartbeat_interval_seconds: int, 

75) -> None: 

76 """ 

77 Updates the status of workers that have an old last heartbeat time to OFFLINE. 

78 

79 An old heartbeat is one that is more than the worker's heartbeat interval 

80 multiplied by the inactivity_heartbeat_multiple seconds ago. 

81 """ 

82 async with db.session_context(begin_transaction=True) as session: 

83 worker_update_stmt = ( 

84 sa.update(db.Worker) 

85 .values(status=WorkerStatus.OFFLINE) 

86 .where( 

87 sa.func.date_diff_seconds(db.Worker.last_heartbeat_time) 

88 > ( 

89 sa.func.coalesce( 

90 db.Worker.heartbeat_interval_seconds, 

91 sa.bindparam("default_interval", sa.Integer), 

92 ) 

93 * sa.bindparam("multiplier", sa.Integer) 

94 ), 

95 db.Worker.status == WorkerStatus.ONLINE, 

96 ) 

97 ) 

98 

99 result = await session.execute( 

100 worker_update_stmt, 

101 { 

102 "multiplier": inactivity_heartbeat_multiple, 

103 "default_interval": fallback_heartbeat_interval_seconds, 

104 }, 

105 ) 

106 

107 if result.rowcount: 

108 logger.info(f"Marked {result.rowcount} workers as offline.") 

109 

110 

111async def _mark_work_pools_as_not_ready(db: PrefectDBInterface) -> None: 

112 """ 

113 Marks work pools as not ready if they have no online workers. 

114 

115 Emits an event and updates any bookkeeping fields on the work pool. 

116 """ 

117 async with db.session_context(begin_transaction=True) as session: 

118 work_pools_select_stmt = ( 

119 sa.select(db.WorkPool) 

120 .filter(db.WorkPool.status == "READY") 

121 .outerjoin( 

122 db.Worker, 

123 sa.and_( 

124 db.Worker.work_pool_id == db.WorkPool.id, 

125 db.Worker.status == "ONLINE", 

126 ), 

127 ) 

128 .group_by(db.WorkPool.id) 

129 .having(sa.func.count(db.Worker.id) == 0) 

130 ) 

131 

132 result = await session.execute(work_pools_select_stmt) 

133 work_pools = result.scalars().all() 

134 

135 for work_pool in work_pools: 135 ↛ 136line 135 didn't jump to line 136 because the loop on line 135 never started

136 await models.workers.update_work_pool( 

137 session=session, 

138 work_pool_id=work_pool.id, 

139 work_pool=InternalWorkPoolUpdate(status=WorkPoolStatus.NOT_READY), 

140 emit_status_change=emit_work_pool_status_event, 

141 ) 

142 

143 logger.info(f"Marked work pool {work_pool.id} as NOT_READY.") 

144 

145 

146async def _mark_deployments_as_not_ready( 

147 db: PrefectDBInterface, 

148 deployment_last_polled_timeout_seconds: int, 

149) -> None: 

150 """ 

151 Marks deployments as NOT_READY based on their last_polled field. 

152 

153 Emits an event and updates any bookkeeping fields on the deployment. 

154 """ 

155 async with db.session_context(begin_transaction=True) as session: 

156 status_timeout_threshold = now("UTC") - timedelta( 

157 seconds=deployment_last_polled_timeout_seconds 

158 ) 

159 deployment_id_select_stmt = ( 

160 sa.select(db.Deployment.id) 

161 .outerjoin(db.WorkQueue, db.WorkQueue.id == db.Deployment.work_queue_id) 

162 .filter(db.Deployment.status == DeploymentStatus.READY) 

163 .filter(db.Deployment.last_polled.isnot(None)) 

164 .filter( 

165 sa.or_( 

166 # if work_queue.last_polled doesn't exist, use only deployment's 

167 # last_polled 

168 sa.and_( 

169 db.WorkQueue.last_polled.is_(None), 

170 db.Deployment.last_polled < status_timeout_threshold, 

171 ), 

172 # if work_queue.last_polled exists, both times should be less than 

173 # the threshold 

174 sa.and_( 

175 db.WorkQueue.last_polled.isnot(None), 

176 db.Deployment.last_polled < status_timeout_threshold, 

177 db.WorkQueue.last_polled < status_timeout_threshold, 

178 ), 

179 ) 

180 ) 

181 ) 

182 result = await session.execute(deployment_id_select_stmt) 

183 

184 deployment_ids_to_mark_unready = result.scalars().all() 

185 

186 await mark_deployments_not_ready( 

187 deployment_ids=deployment_ids_to_mark_unready, 

188 ) 

189 

190 

191async def _mark_work_queues_as_not_ready( 

192 db: PrefectDBInterface, 

193 work_queue_last_polled_timeout_seconds: int, 

194) -> None: 

195 """ 

196 Marks work queues as NOT_READY based on their last_polled field. 

197 """ 

198 async with db.session_context(begin_transaction=True) as session: 

199 status_timeout_threshold = now("UTC") - timedelta( 

200 seconds=work_queue_last_polled_timeout_seconds 

201 ) 

202 id_select_stmt = ( 

203 sa.select(db.WorkQueue.id) 

204 .outerjoin(db.WorkPool, db.WorkPool.id == db.WorkQueue.work_pool_id) 

205 .filter(db.WorkQueue.status == "READY") 

206 .filter(db.WorkQueue.last_polled.isnot(None)) 

207 .filter(db.WorkQueue.last_polled < status_timeout_threshold) 

208 .order_by(db.WorkQueue.last_polled.asc()) 

209 ) 

210 result = await session.execute(id_select_stmt) 

211 unready_work_queue_ids = result.scalars().all() 

212 

213 await mark_work_queues_not_ready( 

214 work_queue_ids=unready_work_queue_ids, 

215 )