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

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"]`): 

5 

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. 

9 

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. 

15 

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`. 

19 

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

24 

25from __future__ import annotations 

26 

27import asyncio 

28import logging 

29from contextlib import asynccontextmanager 

30from datetime import timedelta 

31from typing import AsyncIterator 

32 

33import sqlalchemy as sa 

34from docket import CurrentDocket, Depends, Docket, Perpetual 

35from sqlalchemy.ext.asyncio import AsyncSession 

36 

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 

44 

45logger: logging.Logger = get_logger(__name__) 

46 

47 

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] = {} 

57 

58 

59def _maintenance_database_config( 

60 db: PrefectDBInterface, 

61) -> AsyncPostgresConfiguration | None: 

62 """Return a Postgres config with no statement timeout for vacuum work. 

63 

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 

80 

81 

82@asynccontextmanager 

83async def _maintenance_session( 

84 db: PrefectDBInterface, 

85) -> AsyncIterator[AsyncSession]: 

86 """A transactional session for vacuum maintenance queries. 

87 

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 

101 

102 

103# --------------------------------------------------------------------------- 

104# Finder (perpetual service) 

105# --------------------------------------------------------------------------- 

106 

107 

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. 

124 

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. 

128 

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

138 

139 

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. 

158 

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

169 

170 

171# --------------------------------------------------------------------------- 

172# Cleanup tasks (docket task functions) 

173# --------------------------------------------------------------------------- 

174 

175 

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) 

195 

196 

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) 

216 

217 

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. 

223 

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 ) 

237 

238 

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) 

259 

260 

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. 

266 

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 

275 

276 for event_type, type_retention in overrides.items(): 

277 retention = min(type_retention, global_retention) 

278 retention_cutoff = now("UTC") - retention 

279 

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 ) 

295 

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 ) 

313 

314 

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 

323 

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 ) 

335 

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 ) 

349 

350 

351# --------------------------------------------------------------------------- 

352# Helpers 

353# --------------------------------------------------------------------------- 

354 

355 

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. 

361 

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. 

365 

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 ) 

375 

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() 

385 

386 if not rows: 

387 break 

388 

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() 

398 

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 

423 

424 await asyncio.sleep(0) 

425 

426 return total_updated, total_deleted 

427 

428 

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