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
« 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"""
5from __future__ import annotations
7import logging
8from datetime import datetime, timedelta, timezone
9from typing import Annotated
10from uuid import UUID
12from docket import CurrentDocket, Depends, Docket, Logged, Perpetual
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
24logger: logging.Logger = get_logger(__name__)
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
41 occupancy_seconds = (
42 datetime.now(timezone.utc) - expired_lease.created_at
43 ).total_seconds()
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 )
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)
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()
76 if expired_lease_ids:
77 logger.info(f"Scheduling revocation of {len(expired_lease_ids)} expired leases")
79 for lease_id in expired_lease_ids:
80 await docket.add(revoke_expired_lease)(lease_id)