Coverage for /usr/local/lib/python3.12/site-packages/prefect/server/services/repossessor.py: 92%

30 statements  

« prev     ^ index     » next       coverage.py v7.15.2, created at 2026-10-07 02:04 +0000

1""" 

2The Repossessor service. Handles reconciliation of expired concurrency leases. 

3""" 

4 

5from __future__ import annotations 

6 

7import logging 

8from datetime import datetime, timedelta, timezone 

9from typing import Annotated 

10from uuid import UUID 

11 

12from docket import CurrentDocket, Depends, Docket, Logged, Perpetual 

13 

14from prefect.logging import get_logger 

15from prefect.server.concurrency.lease_storage import ( 

16 ConcurrencyLeaseStorage, 

17 get_concurrency_lease_storage, 

18) 

19from prefect.server.database import PrefectDBInterface, provide_database_interface 

20from prefect.server.models.concurrency_limits_v2 import bulk_decrement_active_slots 

21from prefect.server.services.perpetual_services import perpetual_service 

22from prefect.settings.context import get_current_settings 

23 

24logger: logging.Logger = get_logger(__name__) 

25 

26 

27async def revoke_expired_lease( 

28 lease_id: Annotated[UUID, Logged], 

29 *, 

30 db: PrefectDBInterface = Depends(provide_database_interface), 

31 lease_storage: ConcurrencyLeaseStorage = Depends(get_concurrency_lease_storage), 

32) -> None: 

33 """Revoke a single expired lease (docket task).""" 

34 expired_lease = await lease_storage.read_lease(lease_id) 

35 if expired_lease is None or expired_lease.metadata is None: 35 ↛ 36line 35 didn't jump to line 36 because the condition on line 35 was never true

36 logger.warning( 

37 f"Lease {lease_id} should be revoked but was not found or has no metadata" 

38 ) 

39 return 

40 

41 occupancy_seconds = ( 

42 datetime.now(timezone.utc) - expired_lease.created_at 

43 ).total_seconds() 

44 

45 logger.info( 

46 f"Revoking lease {lease_id} for {len(expired_lease.resource_ids)} " 

47 f"concurrency limits with {expired_lease.metadata.slots} slots" 

48 ) 

49 

50 async with db.session_context(begin_transaction=True) as session: 

51 await bulk_decrement_active_slots( 

52 session=session, 

53 concurrency_limit_ids=expired_lease.resource_ids, 

54 slots=expired_lease.metadata.slots, 

55 occupancy_seconds=occupancy_seconds, 

56 ) 

57 await lease_storage.revoke_lease(lease_id) 

58 

59 

60@perpetual_service( 

61 enabled_getter=lambda: get_current_settings().server.services.repossessor.enabled, 

62) 

63async def monitor_expired_leases( 

64 docket: Docket = CurrentDocket(), 

65 lease_storage: ConcurrencyLeaseStorage = Depends(get_concurrency_lease_storage), 

66 perpetual: Perpetual = Perpetual( 

67 automatic=True, 

68 every=timedelta( 

69 seconds=get_current_settings().server.services.repossessor.loop_seconds 

70 ), 

71 ), 

72) -> None: 

73 """Monitor for expired leases and schedule revocation tasks.""" 

74 expired_lease_ids = await lease_storage.read_expired_lease_ids() 

75 

76 if expired_lease_ids: 

77 logger.info(f"Scheduling revocation of {len(expired_lease_ids)} expired leases") 

78 

79 for lease_id in expired_lease_ids: 

80 await docket.add(revoke_expired_lease)(lease_id)