Coverage for /usr/local/lib/python3.12/site-packages/prefect/server/services/cleanup_reconciler.py: 93%
24 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 cleanup reconciler service. Handles cleanup message lease expiry.
3"""
5from __future__ import annotations
7import logging
8from datetime import timedelta
10from docket import Depends, Perpetual
12from prefect.logging import get_logger
13from prefect.server.services.perpetual_services import perpetual_service
14from prefect.server.worker_communication.cleanup_queue import (
15 CleanupQueueLeaseExpiryResult,
16 WorkerCleanupQueue,
17 get_worker_cleanup_queue,
18)
19from prefect.settings.context import get_current_settings
21logger: logging.Logger = get_logger(__name__)
23_service_cleanup_queue: WorkerCleanupQueue | None = None
24_service_cleanup_queue_storage: str | None = None
27def _get_service_worker_cleanup_queue() -> WorkerCleanupQueue:
28 global _service_cleanup_queue, _service_cleanup_queue_storage
30 storage = get_current_settings().server.worker_channel.cleanup_queue_storage
31 if _service_cleanup_queue is None or _service_cleanup_queue_storage != storage:
32 _service_cleanup_queue = get_worker_cleanup_queue()
33 _service_cleanup_queue_storage = storage
35 return _service_cleanup_queue
38@perpetual_service(
39 enabled_getter=lambda: (
40 get_current_settings().server.services.cleanup_reconciler.enabled
41 ),
42)
43async def reconcile_cleanup_delivery(
44 cleanup_queue: WorkerCleanupQueue = Depends(_get_service_worker_cleanup_queue),
45 perpetual: Perpetual = Perpetual(
46 automatic=True,
47 every=timedelta(
48 seconds=get_current_settings().server.services.cleanup_reconciler.loop_seconds
49 ),
50 ),
51) -> CleanupQueueLeaseExpiryResult:
52 """
53 Reconcile overdue cleanup reservations and apply retry/DLQ policy.
54 """
55 settings = get_current_settings().server.services.cleanup_reconciler
56 result = await cleanup_queue.expire_leases(limit=settings.batch_size)
58 if result.redelivered or result.dead_lettered: 58 ↛ 59line 58 didn't jump to line 59 because the condition on line 58 was never true
59 logger.info(
60 "Expired worker cleanup leases: redelivered=%s dead_lettered=%s",
61 len(result.redelivered),
62 len(result.dead_lettered),
63 )
65 return result