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

1from __future__ import annotations 

2 

3from collections.abc import Iterable, Mapping 

4import importlib 

5import logging 

6from typing import Any, ClassVar, Literal 

7from uuid import UUID 

8 

9from pydantic import ConfigDict, Field 

10 

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 

16 

17try: 

18 from prometheus_client import Counter 

19except ImportError: # pragma: no cover - prometheus_client is a runtime dependency 

20 Counter = None 

21 

22CleanupQueueOperation = Literal["ack", "release", "renew"] 

23 

24logger: logging.Logger = get_logger(__name__) 

25 

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 

46 

47 

48class CleanupQueueMessage(PrefectBaseModel): 

49 model_config: ClassVar[ConfigDict] = ConfigDict(extra="forbid") 

50 

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 

61 

62 

63class CleanupQueueReservation(CleanupQueueMessage): 

64 reservation_token: str = Field(min_length=1) 

65 lease_expires_at: DateTime 

66 delivery_count: PositiveInteger 

67 

68 

69class CleanupQueueDeadLetter(PrefectBaseModel): 

70 model_config: ClassVar[ConfigDict] = ConfigDict(extra="forbid") 

71 

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 

79 

80 

81class CleanupQueueOperationResult(PrefectBaseModel): 

82 model_config: ClassVar[ConfigDict] = ConfigDict(extra="forbid") 

83 

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 

90 

91 

92class CleanupQueueLeaseExpiryResult(PrefectBaseModel): 

93 model_config: ClassVar[ConfigDict] = ConfigDict(extra="forbid") 

94 

95 redelivered: list[CleanupQueueMessage] = Field(default_factory=list) 

96 dead_lettered: list[CleanupQueueDeadLetter] = Field(default_factory=list) 

97 

98 

99class CleanupQueueWakeup(PrefectBaseModel): 

100 model_config: ClassVar[ConfigDict] = ConfigDict(extra="forbid") 

101 

102 work_pool_id: UUID 

103 sequence: PositiveInteger 

104 

105 

106class WorkerCleanupQueue: 

107 """ 

108 Interface for cleanup delivery queue storage. 

109 

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 """ 

116 

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. 

130 

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 ... 

137 

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. 

148 

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 ... 

156 

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. 

166 

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 ... 

173 

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. 

184 

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 ... 

190 

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. 

200 

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 ... 

206 

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. 

215 

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 ... 

221 

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. 

230 

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 ... 

235 

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. 

244 

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 ... 

249 

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. 

253 

254 Implementations should advance and return a monotonic wakeup sequence so 

255 local dispatchers can avoid missing notifications. 

256 """ 

257 ... 

258 

259 async def read_wakeup_sequence(self, work_pool_id: UUID) -> int: 

260 """ 

261 Return the latest wakeup sequence observed for a work pool. 

262 

263 Callers use this value as the `after` cursor when waiting for future 

264 wakeups. 

265 """ 

266 ... 

267 

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`. 

277 

278 Returns the next wakeup when one is observed, or `None` when `timeout` 

279 elapses before a newer wakeup is available. 

280 """ 

281 ... 

282 

283 

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() 

296 

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 ) 

311 

312 

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 

321 

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 ) 

330 

331 

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() 

345 

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 ) 

355 

356 

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() 

375 

376 

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]