Coverage for /usr/local/lib/python3.12/site-packages/prefect/server/services/db_vacuum.py: 59%
124 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 database vacuum service. Two perpetual services schedule cleanup tasks
3independently, gated by the `enabled` set in
4`PREFECT_SERVER_SERVICES_DB_VACUUM_ENABLED` (default `["events"]`):
61. schedule_vacuum_tasks — Cleans up old flow runs and orphaned resources
7 (logs, artifacts, artifact collections). Enabled when `"flow_runs"`
8 is in the enabled set.
102. schedule_event_vacuum_tasks — Cleans up old events, including any
11 event types with per-type retention overrides. Enabled when `"events"`
12 is in the enabled set **and** `event_persister.enabled` is true
13 (the default), so that operators who disabled event processing are not
14 surprised on upgrade. Runs in all server modes, including ephemeral.
16Per-event-type retention can be customised via
17`PREFECT_SERVER_SERVICES_DB_VACUUM_EVENT_RETENTION_OVERRIDES`. Event types
18not listed fall back to `server.events.retention_period`.
20Each task runs independently with its own error isolation and
21docket-managed retries. Deterministic keys prevent duplicate tasks from
22accumulating if a cycle overlaps with in-progress work.
23"""
25from __future__ import annotations
27import asyncio
28import logging
29from contextlib import asynccontextmanager
30from datetime import timedelta
31from typing import AsyncIterator
33import sqlalchemy as sa
34from docket import CurrentDocket, Depends, Docket, Perpetual
35from sqlalchemy.ext.asyncio import AsyncSession
37from prefect.logging import get_logger
38from prefect.server.database import PrefectDBInterface, provide_database_interface
39from prefect.server.database.configurations import AsyncPostgresConfiguration
40from prefect.server.schemas.states import TERMINAL_STATES
41from prefect.server.services.perpetual_services import perpetual_service
42from prefect.settings.context import get_current_settings
43from prefect.types._datetime import now
45logger: logging.Logger = get_logger(__name__)
48# Vacuum runs batched maintenance deletes that legitimately scan large tables
49# (e.g. the orphaned-log anti-join). On Postgres these inherit the asyncpg
50# `command_timeout` derived from `PREFECT_API_DATABASE_TIMEOUT` (10s by
51# default) — a latency budget meant for user-facing API queries, not bulk
52# maintenance. When a batch exceeds it asyncpg raises `TimeoutError`, killing
53# the task before it makes progress. Maintenance work runs on a dedicated
54# connection with no statement timeout so it can run to completion; sqlite's
55# `timeout` is a lock-wait, not a statement deadline, so it is left as-is.
56_MAINTENANCE_CONFIGS: dict[str, AsyncPostgresConfiguration] = {}
59def _maintenance_database_config(
60 db: PrefectDBInterface,
61) -> AsyncPostgresConfiguration | None:
62 """Return a Postgres config with no statement timeout for vacuum work.
64 Returns `None` for non-Postgres backends, signalling callers to use the
65 default session.
66 """
67 config = db.database_config
68 if not isinstance(config, AsyncPostgresConfiguration): 68 ↛ 69line 68 didn't jump to line 69 because the condition on line 68 was never true
69 return None
70 cached = _MAINTENANCE_CONFIGS.get(config.connection_url)
71 if cached is None:
72 cached = AsyncPostgresConfiguration(connection_url=config.connection_url)
73 # Opt out of the API statement timeout and keep a minimal pool, since
74 # vacuum tasks run sequentially on an hourly loop.
75 cached.timeout = None
76 cached.sqlalchemy_pool_size = 1
77 cached.sqlalchemy_max_overflow = 0
78 _MAINTENANCE_CONFIGS[config.connection_url] = cached
79 return cached
82@asynccontextmanager
83async def _maintenance_session(
84 db: PrefectDBInterface,
85) -> AsyncIterator[AsyncSession]:
86 """A transactional session for vacuum maintenance queries.
88 On Postgres this uses a dedicated connection with no statement timeout;
89 other backends fall back to the default session context.
90 """
91 config = _maintenance_database_config(db)
92 if config is None: 92 ↛ 93line 92 didn't jump to line 93 because the condition on line 92 was never true
93 async with db.session_context(begin_transaction=True) as session:
94 yield session
95 return
96 engine = await config.engine()
97 session = await config.session(engine)
98 async with session:
99 async with config.begin_transaction(session):
100 yield session
103# ---------------------------------------------------------------------------
104# Finder (perpetual service)
105# ---------------------------------------------------------------------------
108@perpetual_service(
109 enabled_getter=lambda: (
110 "flow_runs"
111 in get_current_settings().server.services.db_vacuum.enabled_vacuum_types
112 ),
113)
114async def schedule_vacuum_tasks(
115 docket: Docket = CurrentDocket(),
116 perpetual: Perpetual = Perpetual(
117 automatic=True,
118 every=timedelta(
119 seconds=get_current_settings().server.services.db_vacuum.loop_seconds
120 ),
121 ),
122) -> None:
123 """Schedule cleanup tasks for old flow runs and orphaned resources.
125 Each task is enqueued with a deterministic key so that overlapping
126 cycles (e.g. when cleanup takes longer than loop_seconds) naturally
127 deduplicate instead of piling up redundant work.
129 Disabled by default because it permanently deletes flow runs. Enable
130 via PREFECT_SERVER_SERVICES_DB_VACUUM_ENABLED=true.
131 """
132 await docket.add(vacuum_orphaned_logs, key="db-vacuum:orphaned-logs")()
133 await docket.add(vacuum_orphaned_artifacts, key="db-vacuum:orphaned-artifacts")()
134 await docket.add(
135 vacuum_stale_artifact_collections, key="db-vacuum:stale-collections"
136 )()
137 await docket.add(vacuum_old_flow_runs, key="db-vacuum:old-flow-runs")()
140@perpetual_service(
141 enabled_getter=lambda: (
142 "events"
143 in get_current_settings().server.services.db_vacuum.enabled_vacuum_types
144 and get_current_settings().server.services.event_persister.enabled
145 ),
146 run_in_ephemeral=True,
147)
148async def schedule_event_vacuum_tasks(
149 docket: Docket = CurrentDocket(),
150 perpetual: Perpetual = Perpetual(
151 automatic=True,
152 every=timedelta(
153 seconds=get_current_settings().server.services.db_vacuum.loop_seconds
154 ),
155 ),
156) -> None:
157 """Schedule cleanup tasks for old events and heartbeat events.
159 Enabled by default (`"events"` is in the default enabled set).
160 Automatically disabled when the event persister service is disabled
161 (PREFECT_SERVER_SERVICES_EVENT_PERSISTER_ENABLED=false) so that
162 operators who opted out of event processing are not surprised by
163 trimming on upgrade.
164 """
165 await docket.add(
166 vacuum_events_with_retention_overrides, key="db-vacuum:retention-overrides"
167 )()
168 await docket.add(vacuum_old_events, key="db-vacuum:old-events")()
171# ---------------------------------------------------------------------------
172# Cleanup tasks (docket task functions)
173# ---------------------------------------------------------------------------
176async def vacuum_orphaned_logs(
177 *,
178 db: PrefectDBInterface = Depends(provide_database_interface),
179) -> None:
180 """Delete logs whose flow_run_id references a non-existent flow run."""
181 settings = get_current_settings().server.services.db_vacuum
182 deleted = await _batch_delete(
183 db,
184 db.Log,
185 sa.and_(
186 db.Log.flow_run_id.is_not(None),
187 ~sa.exists(
188 sa.select(sa.literal(1)).where(db.FlowRun.id == db.Log.flow_run_id)
189 ),
190 ),
191 settings.batch_size,
192 )
193 if deleted:
194 logger.info("Database vacuum: deleted %d orphaned logs.", deleted)
197async def vacuum_orphaned_artifacts(
198 *,
199 db: PrefectDBInterface = Depends(provide_database_interface),
200) -> None:
201 """Delete artifacts whose flow_run_id references a non-existent flow run."""
202 settings = get_current_settings().server.services.db_vacuum
203 deleted = await _batch_delete(
204 db,
205 db.Artifact,
206 sa.and_(
207 db.Artifact.flow_run_id.is_not(None),
208 ~sa.exists(
209 sa.select(sa.literal(1)).where(db.FlowRun.id == db.Artifact.flow_run_id)
210 ),
211 ),
212 settings.batch_size,
213 )
214 if deleted:
215 logger.info("Database vacuum: deleted %d orphaned artifacts.", deleted)
218async def vacuum_stale_artifact_collections(
219 *,
220 db: PrefectDBInterface = Depends(provide_database_interface),
221) -> None:
222 """Reconcile artifact collections whose latest_id points to a deleted artifact.
224 Re-points to the next latest version if one exists, otherwise deletes
225 the collection row.
226 """
227 settings = get_current_settings().server.services.db_vacuum
228 updated, deleted = await _reconcile_artifact_collections(db, settings.batch_size)
229 if updated or deleted:
230 logger.info(
231 "Database vacuum: reconciled %d stale artifact collections "
232 "(%d re-pointed, %d removed).",
233 updated + deleted,
234 updated,
235 deleted,
236 )
239async def vacuum_old_flow_runs(
240 *,
241 db: PrefectDBInterface = Depends(provide_database_interface),
242) -> None:
243 """Delete old top-level terminal flow runs past the retention period."""
244 settings = get_current_settings().server.services.db_vacuum
245 retention_cutoff = now("UTC") - settings.retention_period
246 deleted = await _batch_delete(
247 db,
248 db.FlowRun,
249 sa.and_(
250 db.FlowRun.parent_task_run_id.is_(None),
251 db.FlowRun.state_type.in_(TERMINAL_STATES),
252 db.FlowRun.end_time.is_not(None),
253 db.FlowRun.end_time < retention_cutoff,
254 ),
255 settings.batch_size,
256 )
257 if deleted:
258 logger.info("Database vacuum: deleted %d old flow runs.", deleted)
261async def vacuum_events_with_retention_overrides(
262 *,
263 db: PrefectDBInterface = Depends(provide_database_interface),
264) -> None:
265 """Delete events whose types have per-type retention overrides.
267 Iterates over all entries in `event_retention_overrides` and deletes
268 events (and their resources) that are older than the configured retention
269 for that type, capped by the global events retention period.
270 """
271 settings = get_current_settings()
272 global_retention = settings.server.events.retention_period
273 overrides = settings.server.services.db_vacuum.event_retention_overrides
274 batch_size = settings.server.services.db_vacuum.batch_size
276 for event_type, type_retention in overrides.items():
277 retention = min(type_retention, global_retention)
278 retention_cutoff = now("UTC") - retention
280 # Delete event resources first (no FK cascade)
281 event_ids = (
282 sa.select(db.Event.id)
283 .where(
284 db.Event.event == event_type,
285 db.Event.occurred < retention_cutoff,
286 )
287 .scalar_subquery()
288 )
289 resources_deleted = await _batch_delete(
290 db,
291 db.EventResource,
292 db.EventResource.event_id.in_(event_ids),
293 batch_size,
294 )
296 # Then delete the events themselves
297 events_deleted = await _batch_delete(
298 db,
299 db.Event,
300 sa.and_(
301 db.Event.event == event_type,
302 db.Event.occurred < retention_cutoff,
303 ),
304 batch_size,
305 )
306 if events_deleted or resources_deleted: 306 ↛ 307line 306 didn't jump to line 307 because the condition on line 306 was never true
307 logger.info(
308 "Database vacuum: deleted %d %r events and %d event resources.",
309 events_deleted,
310 event_type,
311 resources_deleted,
312 )
315async def vacuum_old_events(
316 *,
317 db: PrefectDBInterface = Depends(provide_database_interface),
318) -> None:
319 """Delete all events and event resources past the general events retention period."""
320 settings = get_current_settings()
321 retention_cutoff = now("UTC") - settings.server.events.retention_period
322 batch_size = settings.server.services.db_vacuum.batch_size
324 # Delete old event resources first (no FK cascade on event_id).
325 # Uses EventResource.occurred (the event timestamp) rather than
326 # EventResource.updated (the row insertion time) so that retention
327 # is measured from when the event happened, consistent with how
328 # events themselves are deleted by Event.occurred below.
329 resources_deleted = await _batch_delete(
330 db,
331 db.EventResource,
332 db.EventResource.occurred < retention_cutoff,
333 batch_size,
334 )
336 # Then delete old events
337 events_deleted = await _batch_delete(
338 db,
339 db.Event,
340 db.Event.occurred < retention_cutoff,
341 batch_size,
342 )
343 if events_deleted or resources_deleted: 343 ↛ 344line 343 didn't jump to line 344 because the condition on line 343 was never true
344 logger.info(
345 "Database vacuum: deleted %d old events and %d event resources.",
346 events_deleted,
347 resources_deleted,
348 )
351# ---------------------------------------------------------------------------
352# Helpers
353# ---------------------------------------------------------------------------
356async def _reconcile_artifact_collections(
357 db: PrefectDBInterface,
358 batch_size: int,
359) -> tuple[int, int]:
360 """Reconcile artifact collections whose latest_id points to a deleted artifact.
362 For each stale collection, if another artifact with the same key still
363 exists, re-point latest_id to the newest remaining version (mirroring the
364 logic in models.artifacts.delete_artifact). Otherwise delete the row.
366 Returns (updated_count, deleted_count).
367 """
368 total_updated = 0
369 total_deleted = 0
370 stale_condition = ~sa.exists(
371 sa.select(sa.literal(1)).where(
372 db.Artifact.id == db.ArtifactCollection.latest_id
373 )
374 )
376 while True:
377 async with _maintenance_session(db) as session:
378 rows = (
379 await session.execute(
380 sa.select(db.ArtifactCollection.id, db.ArtifactCollection.key)
381 .where(stale_condition)
382 .limit(batch_size)
383 )
384 ).all()
386 if not rows:
387 break
389 for collection_id, key in rows:
390 next_latest = (
391 await session.execute(
392 sa.select(db.Artifact)
393 .where(db.Artifact.key == key)
394 .order_by(db.Artifact.created.desc())
395 .limit(1)
396 )
397 ).scalar_one_or_none()
399 if next_latest is not None:
400 await session.execute(
401 sa.update(db.ArtifactCollection)
402 .where(db.ArtifactCollection.id == collection_id)
403 .values(
404 latest_id=next_latest.id,
405 data=next_latest.data,
406 description=next_latest.description,
407 type=next_latest.type,
408 created=next_latest.created,
409 updated=next_latest.updated,
410 flow_run_id=next_latest.flow_run_id,
411 task_run_id=next_latest.task_run_id,
412 metadata_=next_latest.metadata_,
413 )
414 )
415 total_updated += 1
416 else:
417 await session.execute(
418 sa.delete(db.ArtifactCollection).where(
419 db.ArtifactCollection.id == collection_id
420 )
421 )
422 total_deleted += 1
424 await asyncio.sleep(0)
426 return total_updated, total_deleted
429async def _batch_delete(
430 db: PrefectDBInterface,
431 model: type,
432 condition: sa.ColumnElement[bool],
433 batch_size: int,
434) -> int:
435 """Delete matching rows in batches. Each batch gets its own DB transaction."""
436 total = 0
437 while True:
438 async with _maintenance_session(db) as session:
439 subquery = (
440 sa.select(model.id).where(condition).limit(batch_size).scalar_subquery()
441 )
442 result = await session.execute(
443 sa.delete(model).where(model.id.in_(subquery))
444 )
445 deleted = result.rowcount
446 if deleted == 0: 446 ↛ 448line 446 didn't jump to line 448 because the condition on line 446 was always true
447 break
448 total += deleted
449 await asyncio.sleep(0) # yield to event loop between batches
450 return total