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

1""" 

2The cleanup reconciler service. Handles cleanup message lease expiry. 

3""" 

4 

5from __future__ import annotations 

6 

7import logging 

8from datetime import timedelta 

9 

10from docket import Depends, Perpetual 

11 

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 

20 

21logger: logging.Logger = get_logger(__name__) 

22 

23_service_cleanup_queue: WorkerCleanupQueue | None = None 

24_service_cleanup_queue_storage: str | None = None 

25 

26 

27def _get_service_worker_cleanup_queue() -> WorkerCleanupQueue: 

28 global _service_cleanup_queue, _service_cleanup_queue_storage 

29 

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 

34 

35 return _service_cleanup_queue 

36 

37 

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) 

57 

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 ) 

64 

65 return result