Coverage for /usr/local/lib/python3.12/site-packages/prefect/server/worker_communication/cleanup_queue/__init__.py: 72%
114 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
1from __future__ import annotations
3from collections.abc import Iterable, Mapping
4import importlib
5import logging
6from typing import Any, ClassVar, Literal
7from uuid import UUID
9from pydantic import ConfigDict, Field
11from prefect._internal.schemas.bases import PrefectBaseModel
12from prefect.client.schemas.worker_channel import CleanupKind, CleanupOperationStatus
13from prefect.logging import get_logger
14from prefect.settings.context import get_current_settings
15from prefect.types import DateTime, NonNegativeInteger, PositiveInteger
17try:
18 from prometheus_client import Counter
19except ImportError: # pragma: no cover - prometheus_client is a runtime dependency
20 Counter = None
22CleanupQueueOperation = Literal["ack", "release", "renew"]
24logger: logging.Logger = get_logger(__name__)
26if Counter is not None: 26 ↛ 43line 26 didn't jump to line 43 because the condition on line 26 was always true
27 CLEANUP_QUEUE_DEAD_LETTERS = Counter(
28 "prefect_worker_cleanup_queue_dead_letters_total",
29 "Worker cleanup queue messages moved to the dead-letter queue.",
30 ["kind", "reason"],
31 )
32 CLEANUP_QUEUE_LEASE_EXPIRATIONS = Counter(
33 "prefect_worker_cleanup_queue_lease_expirations_total",
34 "Worker cleanup queue lease expirations by transition result.",
35 ["result"],
36 )
37 CLEANUP_QUEUE_OPERATIONS = Counter(
38 "prefect_worker_cleanup_queue_operations_total",
39 "Worker cleanup queue operation outcomes.",
40 ["operation", "status"],
41 )
42else:
43 CLEANUP_QUEUE_DEAD_LETTERS = None
44 CLEANUP_QUEUE_LEASE_EXPIRATIONS = None
45 CLEANUP_QUEUE_OPERATIONS = None
48class CleanupQueueMessage(PrefectBaseModel):
49 model_config: ClassVar[ConfigDict] = ConfigDict(extra="forbid")
51 message_id: UUID
52 idempotency_key: str = Field(min_length=1)
53 work_pool_id: UUID
54 work_queue_id: UUID | None = None
55 kind: CleanupKind
56 target: dict[str, Any] = Field(default_factory=dict)
57 data: dict[str, Any] = Field(default_factory=dict)
58 created_at: DateTime
59 updated_at: DateTime
60 delivery_count: NonNegativeInteger = 0
63class CleanupQueueReservation(CleanupQueueMessage):
64 reservation_token: str = Field(min_length=1)
65 lease_expires_at: DateTime
66 delivery_count: PositiveInteger
69class CleanupQueueDeadLetter(PrefectBaseModel):
70 model_config: ClassVar[ConfigDict] = ConfigDict(extra="forbid")
72 message: CleanupQueueMessage
73 reason: str = Field(min_length=1)
74 final_delivery_count: NonNegativeInteger
75 moved_at: DateTime
76 reservation_token: str | None = None
77 lease_expires_at: DateTime | None = None
78 release_reason: str | None = None
81class CleanupQueueOperationResult(PrefectBaseModel):
82 model_config: ClassVar[ConfigDict] = ConfigDict(extra="forbid")
84 message_id: UUID
85 operation: CleanupQueueOperation
86 status: CleanupOperationStatus
87 lease_expires_at: DateTime | None = None
88 reason: str | None = None
89 dead_letter: CleanupQueueDeadLetter | None = None
92class CleanupQueueLeaseExpiryResult(PrefectBaseModel):
93 model_config: ClassVar[ConfigDict] = ConfigDict(extra="forbid")
95 redelivered: list[CleanupQueueMessage] = Field(default_factory=list)
96 dead_lettered: list[CleanupQueueDeadLetter] = Field(default_factory=list)
99class CleanupQueueWakeup(PrefectBaseModel):
100 model_config: ClassVar[ConfigDict] = ConfigDict(extra="forbid")
102 work_pool_id: UUID
103 sequence: PositiveInteger
106class WorkerCleanupQueue:
107 """
108 Interface for cleanup delivery queue storage.
110 Implementations own cleanup message reservation correctness. WebSocket
111 dispatchers may keep process-local routing state, but ack, release, renew,
112 lease expiry, retry accounting, and DLQ transitions must go through this
113 queue. Implementations own the server retry and lease policy and completed
114 idempotency retention semantics.
115 """
117 async def enqueue(
118 self,
119 *,
120 message_id: UUID,
121 idempotency_key: str,
122 work_pool_id: UUID,
123 kind: CleanupKind,
124 target: Mapping[str, Any],
125 data: Mapping[str, Any] | None = None,
126 work_queue_id: UUID | None = None,
127 ) -> CleanupQueueMessage:
128 """
129 Store a cleanup message if it has not already been produced.
131 Implementations should treat `message_id` and `idempotency_key` as stable
132 producer identifiers within a work pool, returning the existing message
133 for repeated enqueue attempts instead of creating duplicates. The
134 optional `work_queue_id` is advisory targeting metadata.
135 """
136 ...
138 async def reserve(
139 self,
140 *,
141 work_pool_id: UUID,
142 cleanup_kinds: Iterable[CleanupKind] | None = None,
143 preferred_work_queue_ids: Iterable[UUID] | None = None,
144 allow_fallback_to_any_queue: bool = True,
145 ) -> CleanupQueueReservation | None:
146 """
147 Atomically reserve one eligible cleanup message for delivery.
149 A successful reservation must increment the committed delivery count,
150 create exactly one active reservation, and return an unguessable token
151 required for follow-up operations. `preferred_work_queue_ids` should be
152 treated as an advisory preference, with pool-wide fallback controlled by
153 `allow_fallback_to_any_queue`.
154 """
155 ...
157 async def ack(
158 self,
159 *,
160 work_pool_id: UUID,
161 message_id: UUID,
162 reservation_token: str,
163 ) -> CleanupQueueOperationResult:
164 """
165 Complete a reserved cleanup message.
167 The operation must validate the work-pool scope and current reservation
168 token atomically before removing the message from active delivery.
169 Implementations should retain completed idempotency state according to
170 their configured retention policy.
171 """
172 ...
174 async def release(
175 self,
176 *,
177 work_pool_id: UUID,
178 message_id: UUID,
179 reservation_token: str,
180 reason: str,
181 ) -> CleanupQueueOperationResult:
182 """
183 Give up the current reservation without completing the cleanup message.
185 The operation must validate the work-pool scope and current reservation
186 token atomically, then either make the message eligible for redelivery or
187 move it to the dead-letter queue when retry policy is exhausted.
188 """
189 ...
191 async def renew(
192 self,
193 *,
194 work_pool_id: UUID,
195 message_id: UUID,
196 reservation_token: str,
197 ) -> CleanupQueueOperationResult:
198 """
199 Extend the lease for the current reservation.
201 The operation must validate the work-pool scope and current reservation
202 token atomically. Renewing a reservation should not increment delivery
203 count because no new delivery has been committed.
204 """
205 ...
207 async def expire_leases(
208 self,
209 *,
210 limit: int = 100,
211 work_pool_id: UUID | None = None,
212 ) -> CleanupQueueLeaseExpiryResult:
213 """
214 Expire overdue reservations in bounded batches.
216 Expired messages should become eligible for redelivery or move to the
217 dead-letter queue according to retry policy. Implementations may scope
218 the sweep to a work pool when `work_pool_id` is provided.
219 """
220 ...
222 async def read_message(
223 self,
224 *,
225 work_pool_id: UUID,
226 message_id: UUID,
227 ) -> CleanupQueueMessage | None:
228 """
229 Read an active cleanup message by work-pool scope and message ID.
231 This is an inspection helper for messages that have not been acked or
232 dead-lettered. It must not return a message from a different work pool.
233 """
234 ...
236 async def read_dead_letter(
237 self,
238 *,
239 work_pool_id: UUID,
240 message_id: UUID,
241 ) -> CleanupQueueDeadLetter | None:
242 """
243 Read a dead-letter entry by work-pool scope and message ID.
245 This is an inspection helper for terminal cleanup failures. It must not
246 return a dead-letter entry from a different work pool.
247 """
248 ...
250 async def wake_dispatchers(self, work_pool_id: UUID) -> CleanupQueueWakeup:
251 """
252 Notify dispatchers that cleanup work may be available for a work pool.
254 Implementations should advance and return a monotonic wakeup sequence so
255 local dispatchers can avoid missing notifications.
256 """
257 ...
259 async def read_wakeup_sequence(self, work_pool_id: UUID) -> int:
260 """
261 Return the latest wakeup sequence observed for a work pool.
263 Callers use this value as the `after` cursor when waiting for future
264 wakeups.
265 """
266 ...
268 async def wait_for_wakeup(
269 self,
270 work_pool_id: UUID,
271 *,
272 after: int = 0,
273 timeout: float | None = None,
274 ) -> CleanupQueueWakeup | None:
275 """
276 Wait for a work-pool wakeup sequence newer than `after`.
278 Returns the next wakeup when one is observed, or `None` when `timeout`
279 elapses before a newer wakeup is available.
280 """
281 ...
284def record_cleanup_queue_dead_letter(
285 dead_letter: CleanupQueueDeadLetter, *, source: str
286) -> None:
287 """
288 Record observability signals for a cleanup message DLQ transition.
289 """
290 message = dead_letter.message
291 if CLEANUP_QUEUE_DEAD_LETTERS is not None:
292 CLEANUP_QUEUE_DEAD_LETTERS.labels(
293 kind=str(message.kind),
294 reason=dead_letter.reason,
295 ).inc()
297 logger.warning(
298 "Worker cleanup message moved to dead-letter queue: "
299 "message_id=%s cleanup_kind=%s work_pool_id=%s "
300 "final_delivery_count=%s reason=%s release_reason=%s "
301 "lease_expires_at=%s source=%s",
302 message.message_id,
303 message.kind,
304 message.work_pool_id,
305 dead_letter.final_delivery_count,
306 dead_letter.reason,
307 dead_letter.release_reason,
308 dead_letter.lease_expires_at,
309 source,
310 )
313def record_cleanup_queue_lease_expiry_result(
314 result: CleanupQueueLeaseExpiryResult,
315) -> None:
316 """
317 Record aggregate lease-expiry transition metrics.
318 """
319 if CLEANUP_QUEUE_LEASE_EXPIRATIONS is None: 319 ↛ 320line 319 didn't jump to line 320 because the condition on line 319 was never true
320 return
322 if result.redelivered: 322 ↛ 323line 322 didn't jump to line 323 because the condition on line 322 was never true
323 CLEANUP_QUEUE_LEASE_EXPIRATIONS.labels(result="redelivered").inc(
324 len(result.redelivered)
325 )
326 if result.dead_lettered: 326 ↛ 327line 326 didn't jump to line 327 because the condition on line 326 was never true
327 CLEANUP_QUEUE_LEASE_EXPIRATIONS.labels(result="dead_lettered").inc(
328 len(result.dead_lettered)
329 )
332def record_cleanup_queue_operation(
333 operation: str,
334 *,
335 status: str,
336 work_pool_id: UUID,
337 message_id: UUID | None = None,
338 cleanup_kind: str | None = None,
339) -> None:
340 """
341 Record observability signals for a cleanup queue operation.
342 """
343 if CLEANUP_QUEUE_OPERATIONS is not None:
344 CLEANUP_QUEUE_OPERATIONS.labels(operation=operation, status=status).inc()
346 logger.debug(
347 "Worker cleanup queue operation: operation=%s status=%s "
348 "work_pool_id=%s message_id=%s cleanup_kind=%s",
349 operation,
350 status,
351 work_pool_id,
352 message_id,
353 cleanup_kind,
354 )
357def get_worker_cleanup_queue() -> WorkerCleanupQueue:
358 """
359 Return a cleanup queue instance from the configured storage module.
360 """
361 worker_channel_settings = get_current_settings().server.worker_channel
362 storage_module = worker_channel_settings.cleanup_queue_storage
363 cleanup_queue_module = importlib.import_module(storage_module)
364 cleanup_queue_class = getattr(cleanup_queue_module, "WorkerCleanupQueue", None)
365 if ( 365 ↛ 370line 365 didn't jump to line 370 because the condition on line 365 was never true
366 not isinstance(cleanup_queue_class, type)
367 or cleanup_queue_class is WorkerCleanupQueue
368 or not issubclass(cleanup_queue_class, WorkerCleanupQueue)
369 ):
370 raise ValueError(
371 f"The module {storage_module} does not contain a concrete "
372 "WorkerCleanupQueue implementation"
373 )
374 return cleanup_queue_class()
377__all__ = [
378 "CleanupQueueDeadLetter",
379 "CleanupQueueLeaseExpiryResult",
380 "CleanupQueueMessage",
381 "CleanupQueueOperation",
382 "CleanupQueueOperationResult",
383 "CleanupQueueReservation",
384 "CleanupQueueWakeup",
385 "WorkerCleanupQueue",
386 "get_worker_cleanup_queue",
387 "record_cleanup_queue_dead_letter",
388 "record_cleanup_queue_lease_expiry_result",
389 "record_cleanup_queue_operation",
390]