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
« 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"""
5from __future__ import annotations
7import logging
8from datetime import timedelta
10import sqlalchemy as sa
11from docket import Perpetual
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
29logger: logging.Logger = get_logger(__name__)
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.
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()
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 )
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.
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 )
99 result = await session.execute(
100 worker_update_stmt,
101 {
102 "multiplier": inactivity_heartbeat_multiple,
103 "default_interval": fallback_heartbeat_interval_seconds,
104 },
105 )
107 if result.rowcount:
108 logger.info(f"Marked {result.rowcount} workers as offline.")
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.
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 )
132 result = await session.execute(work_pools_select_stmt)
133 work_pools = result.scalars().all()
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 )
143 logger.info(f"Marked work pool {work_pool.id} as NOT_READY.")
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.
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)
184 deployment_ids_to_mark_unready = result.scalars().all()
186 await mark_deployments_not_ready(
187 deployment_ids=deployment_ids_to_mark_unready,
188 )
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()
213 await mark_work_queues_not_ready(
214 work_queue_ids=unready_work_queue_ids,
215 )