Coverage for /usr/local/lib/python3.12/site-packages/prefect/server/models/events.py: 65%

164 statements  

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

1from datetime import datetime, timedelta 

2from typing import Any, Dict, List, MutableMapping, Optional, Set, Union 

3from uuid import UUID 

4 

5from cachetools import TTLCache 

6from sqlalchemy.ext.asyncio import AsyncSession 

7 

8from prefect._internal.uuid7 import uuid7 

9from prefect.server import models, schemas 

10from prefect.server.database.orm_models import ( 

11 ORMDeployment, 

12 ORMFlow, 

13 ORMFlowRun, 

14 ORMFlowRunState, 

15 ORMTaskRun, 

16 ORMTaskRunState, 

17 ORMWorkPool, 

18 ORMWorkQueue, 

19) 

20from prefect.server.events.schemas.events import Event 

21from prefect.server.models import deployments 

22from prefect.server.schemas.statuses import DeploymentStatus 

23from prefect.settings import PREFECT_API_EVENTS_RELATED_RESOURCE_CACHE_TTL 

24from prefect.types._datetime import DateTime, now 

25from prefect.utilities.text import truncated_to 

26 

27ResourceData = Dict[str, Dict[str, Any]] 

28RelatedResourceList = List[Dict[str, str]] 

29 

30 

31# Some users use state messages to convey error messages and large results; let's 

32# truncate them so they don't blow out the size of a message 

33TRUNCATE_STATE_MESSAGES_AT = 100_000 

34 

35 

36_flow_run_resource_data_cache: MutableMapping[UUID, ResourceData] = TTLCache( 

37 maxsize=1000, 

38 ttl=PREFECT_API_EVENTS_RELATED_RESOURCE_CACHE_TTL.value().total_seconds(), 

39) 

40 

41 

42async def flow_run_state_change_event( 

43 session: AsyncSession, 

44 occurred: datetime, 

45 flow_run: ORMFlowRun, 

46 initial_state_id: Optional[UUID], 

47 initial_state: Optional[schemas.states.State], 

48 validated_state_id: Optional[UUID], 

49 validated_state: schemas.states.State, 

50) -> Event: 

51 return Event( 

52 occurred=occurred, 

53 event=f"prefect.flow-run.{validated_state.name}", 

54 resource={ 

55 "prefect.resource.id": f"prefect.flow-run.{flow_run.id}", 

56 "prefect.resource.name": flow_run.name, 

57 "prefect.run-count": str(flow_run.run_count), 

58 "prefect.state-message": truncated_to( 

59 TRUNCATE_STATE_MESSAGES_AT, validated_state.message 

60 ), 

61 "prefect.state-name": validated_state.name or "", 

62 "prefect.state-timestamp": ( 

63 validated_state.timestamp.isoformat() 

64 if validated_state.timestamp 

65 else None 

66 ), 

67 "prefect.state-type": validated_state.type.value, 

68 }, 

69 related=await _flow_run_related_resources_from_orm( 

70 session=session, flow_run=flow_run 

71 ), 

72 payload={ 

73 "intended": { 

74 "from": _state_type(initial_state), 

75 "to": _state_type(validated_state), 

76 }, 

77 "initial_state": state_payload(initial_state), 

78 "validated_state": state_payload(validated_state), 

79 }, 

80 # Here we use the state's ID as the ID of the event as well, in order to 

81 # establish the ordering of the state-change events 

82 id=validated_state_id, 

83 follows=initial_state_id if _timing_is_tight(occurred, initial_state) else None, 

84 ) 

85 

86 

87async def _flow_run_related_resources_from_orm( 

88 session: AsyncSession, flow_run: ORMFlowRun 

89) -> RelatedResourceList: 

90 resource_data = _flow_run_resource_data_cache.get(flow_run.id) 

91 if not resource_data: 

92 flow = await models.flows.read_flow(session=session, flow_id=flow_run.flow_id) 

93 deployment: Optional[ORMDeployment] = None 

94 if flow_run.deployment_id: 

95 deployment = await deployments.read_deployment( 

96 session, deployment_id=flow_run.deployment_id 

97 ) 

98 

99 work_queue = None 

100 if flow_run.work_queue_id: 

101 work_queue = await models.work_queues.read_work_queue( 

102 session, work_queue_id=flow_run.work_queue_id 

103 ) 

104 

105 work_pool = work_queue.work_pool if work_queue is not None else None 

106 

107 task_run: Optional[ORMTaskRun] = None 

108 if flow_run.parent_task_run_id: 

109 task_run = await models.task_runs.read_task_run( 

110 session, 

111 task_run_id=flow_run.parent_task_run_id, 

112 ) 

113 

114 resource_data = _as_resource_data( 

115 flow_run, flow, deployment, work_queue, work_pool, task_run 

116 ) 

117 _flow_run_resource_data_cache[flow_run.id] = resource_data 

118 

119 return _resource_data_as_related_resources( 

120 resource_data, 

121 excluded_kinds=["flow-run"], 

122 ) + _provenance_as_related_resources(flow_run.created_by) 

123 

124 

125def _as_resource_data( 

126 flow_run: ORMFlowRun, 

127 flow: Union[ORMFlow, schemas.core.Flow, None], 

128 deployment: Union[ORMDeployment, schemas.responses.DeploymentResponse, None], 

129 work_queue: Union[ORMWorkQueue, schemas.responses.WorkQueueResponse, None], 

130 work_pool: Union[ORMWorkPool, schemas.core.WorkPool, None], 

131 task_run: Union[ORMTaskRun, schemas.core.TaskRun, None] = None, 

132) -> ResourceData: 

133 return { 

134 "flow-run": { 

135 "id": str(flow_run.id), 

136 "name": flow_run.name, 

137 "tags": flow_run.tags if flow_run.tags else [], 

138 "role": "flow-run", 

139 }, 

140 "flow": ( 

141 { 

142 "id": str(flow.id), 

143 "name": flow.name, 

144 "tags": flow.tags if flow.tags else [], 

145 "role": "flow", 

146 } 

147 if flow 

148 else {} 

149 ), 

150 "deployment": ( 

151 { 

152 "id": str(deployment.id), 

153 "name": deployment.name, 

154 "tags": deployment.tags if deployment.tags else [], 

155 "role": "deployment", 

156 } 

157 if deployment 

158 else {} 

159 ), 

160 "work-queue": ( 

161 { 

162 "id": str(work_queue.id), 

163 "name": work_queue.name, 

164 "tags": [], 

165 "role": "work-queue", 

166 } 

167 if work_queue 

168 else {} 

169 ), 

170 "work-pool": ( 

171 { 

172 "id": str(work_pool.id), 

173 "name": work_pool.name, 

174 "tags": [], 

175 "role": "work-pool", 

176 "type": work_pool.type, 

177 } 

178 if work_pool 

179 else {} 

180 ), 

181 "task-run": ( 

182 { 

183 "id": str(task_run.id), 

184 "name": task_run.name, 

185 "tags": task_run.tags if task_run.tags else [], 

186 "role": "task-run", 

187 } 

188 if task_run 

189 else {} 

190 ), 

191 } 

192 

193 

194def _resource_data_as_related_resources( 

195 resource_data: ResourceData, 

196 excluded_kinds: Optional[List[str]] = None, 

197) -> RelatedResourceList: 

198 related = [] 

199 tags: Set[str] = set() 

200 

201 if excluded_kinds is None: 201 ↛ 202line 201 didn't jump to line 202 because the condition on line 201 was never true

202 excluded_kinds = [] 

203 

204 for kind, data in resource_data.items(): 

205 tags |= set(data.get("tags", [])) 

206 

207 if kind in excluded_kinds or not data: 

208 continue 

209 

210 related_resource = { 

211 "prefect.resource.id": f"prefect.{kind}.{data['id']}", 

212 "prefect.resource.role": data["role"], 

213 "prefect.resource.name": data["name"], 

214 } 

215 

216 if kind == "work-pool": 

217 related_resource["prefect.work-pool.type"] = data["type"] 

218 

219 related.append(related_resource) 

220 

221 related += [ 

222 { 

223 "prefect.resource.id": f"prefect.tag.{tag}", 

224 "prefect.resource.role": "tag", 

225 } 

226 for tag in sorted(tags) 

227 ] 

228 

229 return related 

230 

231 

232def _provenance_as_related_resources( 

233 created_by: Optional[schemas.core.CreatedBy], 

234) -> RelatedResourceList: 

235 if not created_by: 235 ↛ 240line 235 didn't jump to line 240 because the condition on line 235 was always true

236 return [] 

237 

238 resource_id: str 

239 

240 if created_by.type == "DEPLOYMENT": 

241 resource_id = f"prefect.deployment.{created_by.id}" 

242 elif created_by.type == "AUTOMATION": 

243 resource_id = f"prefect.automation.{created_by.id}" 

244 else: 

245 return [] 

246 

247 related = { 

248 "prefect.resource.id": resource_id, 

249 "prefect.resource.role": "creator", 

250 } 

251 if created_by.display_value: 

252 related["prefect.resource.name"] = created_by.display_value 

253 return [related] 

254 

255 

256def _state_type( 

257 state: Union[ORMFlowRunState, ORMTaskRunState, Optional[schemas.states.State]], 

258) -> Optional[str]: 

259 return str(state.type.value) if state else None 

260 

261 

262def state_payload(state: Optional[schemas.states.State]) -> Optional[Dict[str, str]]: 

263 """Given a State, return the essential string parts of it for use in an 

264 event payload""" 

265 if not state: 

266 return None 

267 payload: Dict[str, str] = {"type": state.type.value} 

268 if state.name: 268 ↛ 270line 268 didn't jump to line 270 because the condition on line 268 was always true

269 payload["name"] = state.name 

270 if state.message: 

271 payload["message"] = truncated_to(TRUNCATE_STATE_MESSAGES_AT, state.message) 

272 if state.is_paused(): 

273 payload["pause_reschedule"] = str(state.state_details.pause_reschedule).lower() 

274 return payload 

275 

276 

277def _timing_is_tight( 

278 occurred: datetime, 

279 initial_state: Union[ 

280 ORMFlowRunState, ORMTaskRunState, Optional[schemas.states.State] 

281 ], 

282) -> bool: 

283 # Only connect events with event.follows if the timing here is tight, which will 

284 # help us resolve the order of these events if they happen to be delivered out of 

285 # order. If the preceding state change happened a while back, don't worry about 

286 # it because the order is very likely to be unambiguous. 

287 TIGHT_TIMING = timedelta(minutes=5) 

288 if initial_state and initial_state.timestamp: 

289 return bool(-TIGHT_TIMING < (occurred - initial_state.timestamp) < TIGHT_TIMING) 

290 

291 return False 

292 

293 

294async def _deployment_related_resources( 

295 session: AsyncSession, 

296 deployment: ORMDeployment, 

297) -> RelatedResourceList: 

298 """Get related resources (flow, work-queue, work-pool) for a deployment event.""" 

299 related: RelatedResourceList = [] 

300 

301 flow = await models.flows.read_flow(session=session, flow_id=deployment.flow_id) 

302 if flow is not None: 

303 related.append( 

304 { 

305 "prefect.resource.id": f"prefect.flow.{flow.id}", 

306 "prefect.resource.name": flow.name, 

307 "prefect.resource.role": "flow", 

308 } 

309 ) 

310 

311 work_queue = ( 

312 await models.workers.read_work_queue( 

313 session=session, 

314 work_queue_id=deployment.work_queue_id, 

315 ) 

316 if deployment.work_queue_id 

317 else None 

318 ) 

319 if work_queue is not None: 

320 related.append( 

321 { 

322 "prefect.resource.id": f"prefect.work-queue.{work_queue.id}", 

323 "prefect.resource.name": work_queue.name, 

324 "prefect.resource.role": "work-queue", 

325 } 

326 ) 

327 

328 work_pool = ( 

329 await models.workers.read_work_pool( 

330 session=session, 

331 work_pool_id=work_queue.work_pool_id, 

332 ) 

333 if work_queue and work_queue.work_pool_id 

334 else None 

335 ) 

336 if work_pool is not None: 

337 related.append( 

338 { 

339 "prefect.resource.id": f"prefect.work-pool.{work_pool.id}", 

340 "prefect.resource.name": work_pool.name, 

341 "prefect.work-pool.type": work_pool.type, 

342 "prefect.resource.role": "work-pool", 

343 } 

344 ) 

345 

346 if deployment.storage_document_id: 

347 related.append( 

348 { 

349 "prefect.resource.id": ( 

350 f"prefect.block-document.{deployment.storage_document_id}" 

351 ), 

352 "prefect.resource.role": "storage", 

353 } 

354 ) 

355 

356 if deployment.infrastructure_document_id: 

357 related.append( 

358 { 

359 "prefect.resource.id": ( 

360 f"prefect.block-document.{deployment.infrastructure_document_id}" 

361 ), 

362 "prefect.resource.role": "infrastructure", 

363 } 

364 ) 

365 

366 if deployment.concurrency_limit_id: 

367 related.append( 

368 { 

369 "prefect.resource.id": ( 

370 f"prefect.concurrency-limit.{deployment.concurrency_limit_id}" 

371 ), 

372 "prefect.resource.role": "concurrency-limit", 

373 } 

374 ) 

375 

376 return related 

377 

378 

379async def deployment_status_event( 

380 session: AsyncSession, 

381 deployment_id: UUID, 

382 status: DeploymentStatus, 

383 occurred: DateTime, 

384) -> Event: 

385 deployment = await models.deployments.read_deployment( 

386 session=session, deployment_id=deployment_id 

387 ) 

388 assert deployment 

389 

390 return Event( 

391 occurred=occurred, 

392 event=f"prefect.deployment.{status.in_kebab_case()}", 

393 resource={ 

394 "prefect.resource.id": f"prefect.deployment.{deployment.id}", 

395 "prefect.resource.name": f"{deployment.name}", 

396 }, 

397 related=await _deployment_related_resources(session, deployment), 

398 id=uuid7(), 

399 ) 

400 

401 

402async def deployment_created_event( 

403 session: AsyncSession, 

404 deployment: ORMDeployment, 

405 occurred: DateTime, 

406) -> Event: 

407 """Create an event for deployment creation.""" 

408 return Event( 

409 occurred=occurred, 

410 event="prefect.deployment.created", 

411 resource={ 

412 "prefect.resource.id": f"prefect.deployment.{deployment.id}", 

413 "prefect.resource.name": deployment.name, 

414 }, 

415 related=await _deployment_related_resources(session, deployment), 

416 id=uuid7(), 

417 ) 

418 

419 

420async def deployment_updated_event( 

421 session: AsyncSession, 

422 deployment: ORMDeployment, 

423 changed_fields: Dict[str, Dict[str, Any]], 

424 occurred: DateTime, 

425) -> Event: 

426 """Create an event for deployment field updates.""" 

427 return Event( 

428 occurred=occurred, 

429 event="prefect.deployment.updated", 

430 resource={ 

431 "prefect.resource.id": f"prefect.deployment.{deployment.id}", 

432 "prefect.resource.name": deployment.name, 

433 }, 

434 related=await _deployment_related_resources(session, deployment), 

435 payload={ 

436 "updated_fields": list(changed_fields.keys()), 

437 "updates": changed_fields, 

438 }, 

439 id=uuid7(), 

440 ) 

441 

442 

443async def deployment_deleted_event( 

444 session: AsyncSession, 

445 deployment: ORMDeployment, 

446 occurred: DateTime, 

447) -> Event: 

448 """Create an event for deployment deletion.""" 

449 return Event( 

450 occurred=occurred, 

451 event="prefect.deployment.deleted", 

452 resource={ 

453 "prefect.resource.id": f"prefect.deployment.{deployment.id}", 

454 "prefect.resource.name": deployment.name, 

455 }, 

456 related=await _deployment_related_resources(session, deployment), 

457 id=uuid7(), 

458 ) 

459 

460 

461async def work_queue_status_event( 

462 session: AsyncSession, 

463 work_queue: "ORMWorkQueue", 

464 occurred: DateTime, 

465) -> Event: 

466 related_work_pool_info: List[Dict[str, Any]] = [] 

467 

468 if work_queue.work_pool_id: 468 ↛ 484line 468 didn't jump to line 484 because the condition on line 468 was always true

469 work_pool = await models.workers.read_work_pool( 

470 session=session, 

471 work_pool_id=work_queue.work_pool_id, 

472 ) 

473 

474 if work_pool and work_pool.id and work_pool.name: 474 ↛ 484line 474 didn't jump to line 484 because the condition on line 474 was always true

475 related_work_pool_info.append( 

476 { 

477 "prefect.resource.id": f"prefect.work-pool.{work_pool.id}", 

478 "prefect.resource.name": work_pool.name, 

479 "prefect.work-pool.type": work_pool.type, 

480 "prefect.resource.role": "work-pool", 

481 } 

482 ) 

483 

484 return Event( 

485 occurred=occurred, 

486 event=f"prefect.work-queue.{work_queue.status.in_kebab_case()}", 

487 resource={ 

488 "prefect.resource.id": f"prefect.work-queue.{work_queue.id}", 

489 "prefect.resource.name": work_queue.name, 

490 "prefect.resource.role": "work-queue", 

491 }, 

492 related=related_work_pool_info, 

493 id=uuid7(), 

494 ) 

495 

496 

497async def work_pool_status_event( 

498 event_id: UUID, 

499 occurred: DateTime, 

500 pre_update_work_pool: Optional["ORMWorkPool"], 

501 work_pool: "ORMWorkPool", 

502) -> Event: 

503 assert work_pool.status 

504 

505 return Event( 

506 id=event_id, 

507 occurred=occurred, 

508 event=f"prefect.work-pool.{work_pool.status.in_kebab_case()}", 

509 resource={ 

510 "prefect.resource.id": f"prefect.work-pool.{work_pool.id}", 

511 "prefect.resource.name": work_pool.name, 

512 "prefect.work-pool.type": work_pool.type, 

513 }, 

514 follows=_get_recent_preceding_work_pool_event_id(pre_update_work_pool), 

515 ) 

516 

517 

518async def work_pool_updated_event( 

519 session: AsyncSession, 

520 work_pool: "ORMWorkPool", 

521 changed_fields: Dict[ 

522 str, Dict[str, Any] 

523 ], # {"field_name": {"from": value, "to": value}} 

524 occurred: DateTime, 

525) -> Event: 

526 """Create an event for work pool field updates (non-status).""" 

527 return Event( 

528 occurred=occurred, 

529 event="prefect.work-pool.updated", 

530 resource={ 

531 "prefect.resource.id": f"prefect.work-pool.{work_pool.id}", 

532 "prefect.resource.name": work_pool.name, 

533 "prefect.work-pool.type": work_pool.type, 

534 "prefect.resource.role": "work-pool", 

535 }, 

536 payload={ 

537 "updated_fields": list(changed_fields.keys()), 

538 "updates": changed_fields, 

539 }, 

540 id=uuid7(), 

541 ) 

542 

543 

544async def work_queue_updated_event( 

545 session: AsyncSession, 

546 work_queue: "ORMWorkQueue", 

547 changed_fields: Dict[str, Dict[str, Any]], 

548 occurred: DateTime, 

549) -> Event: 

550 """Create an event for work queue field updates (non-status).""" 

551 related_work_pool_info: List[Dict[str, Any]] = [] 

552 

553 if work_queue.work_pool_id: 553 ↛ 568line 553 didn't jump to line 568 because the condition on line 553 was always true

554 work_pool = await models.workers.read_work_pool( 

555 session=session, 

556 work_pool_id=work_queue.work_pool_id, 

557 ) 

558 if work_pool and work_pool.id and work_pool.name: 

559 related_work_pool_info.append( 

560 { 

561 "prefect.resource.id": f"prefect.work-pool.{work_pool.id}", 

562 "prefect.resource.name": work_pool.name, 

563 "prefect.work-pool.type": work_pool.type, 

564 "prefect.resource.role": "work-pool", 

565 } 

566 ) 

567 

568 return Event( 

569 occurred=occurred, 

570 event="prefect.work-queue.updated", 

571 resource={ 

572 "prefect.resource.id": f"prefect.work-queue.{work_queue.id}", 

573 "prefect.resource.name": work_queue.name, 

574 "prefect.resource.role": "work-queue", 

575 }, 

576 related=related_work_pool_info, 

577 payload={ 

578 "updated_fields": list(changed_fields.keys()), 

579 "updates": changed_fields, 

580 }, 

581 id=uuid7(), 

582 ) 

583 

584 

585def _work_pool_related_resources(work_pool: "ORMWorkPool") -> RelatedResourceList: 

586 """The work pool's default queue and result-storage block as related 

587 resources, each with a role naming its relationship to the pool.""" 

588 related: RelatedResourceList = [] 

589 

590 if work_pool.default_queue_id: 590 ↛ 600line 590 didn't jump to line 600 because the condition on line 590 was always true

591 related.append( 

592 { 

593 "prefect.resource.id": ( 

594 f"prefect.work-queue.{work_pool.default_queue_id}" 

595 ), 

596 "prefect.resource.role": "default-queue", 

597 } 

598 ) 

599 

600 storage_configuration = work_pool.storage_configuration 

601 result_storage_block_id = getattr( 

602 storage_configuration, "default_result_storage_block_id", None 

603 ) 

604 if result_storage_block_id: 

605 related.append( 

606 { 

607 "prefect.resource.id": ( 

608 f"prefect.block-document.{result_storage_block_id}" 

609 ), 

610 "prefect.resource.role": "result-storage", 

611 } 

612 ) 

613 

614 return related 

615 

616 

617def _work_pool_resource(work_pool: "ORMWorkPool") -> Dict[str, str]: 

618 return { 

619 "prefect.resource.id": f"prefect.work-pool.{work_pool.id}", 

620 "prefect.resource.name": work_pool.name, 

621 "prefect.work-pool.type": work_pool.type, 

622 "prefect.resource.role": "work-pool", 

623 } 

624 

625 

626async def work_pool_created_event( 

627 work_pool: "ORMWorkPool", 

628 occurred: DateTime, 

629) -> Event: 

630 """Create an event for work pool creation.""" 

631 return Event( 

632 occurred=occurred, 

633 event="prefect.work-pool.created", 

634 resource=_work_pool_resource(work_pool), 

635 related=_work_pool_related_resources(work_pool), 

636 id=uuid7(), 

637 ) 

638 

639 

640async def work_pool_deleted_event( 

641 work_pool: "ORMWorkPool", 

642 occurred: DateTime, 

643) -> Event: 

644 """Create an event for work pool deletion.""" 

645 return Event( 

646 occurred=occurred, 

647 event="prefect.work-pool.deleted", 

648 resource=_work_pool_resource(work_pool), 

649 related=_work_pool_related_resources(work_pool), 

650 id=uuid7(), 

651 ) 

652 

653 

654async def _work_queue_work_pool_related( 

655 session: AsyncSession, 

656 work_queue: "ORMWorkQueue", 

657) -> List[Dict[str, Any]]: 

658 related: List[Dict[str, Any]] = [] 

659 if work_queue.work_pool_id: 659 ↛ 673line 659 didn't jump to line 673 because the condition on line 659 was always true

660 work_pool = await models.workers.read_work_pool( 

661 session=session, 

662 work_pool_id=work_queue.work_pool_id, 

663 ) 

664 if work_pool and work_pool.id and work_pool.name: 

665 related.append( 

666 { 

667 "prefect.resource.id": f"prefect.work-pool.{work_pool.id}", 

668 "prefect.resource.name": work_pool.name, 

669 "prefect.work-pool.type": work_pool.type, 

670 "prefect.resource.role": "work-pool", 

671 } 

672 ) 

673 return related 

674 

675 

676async def work_queue_created_event( 

677 session: AsyncSession, 

678 work_queue: "ORMWorkQueue", 

679 occurred: DateTime, 

680) -> Event: 

681 """Create an event for work queue creation.""" 

682 return Event( 

683 occurred=occurred, 

684 event="prefect.work-queue.created", 

685 resource={ 

686 "prefect.resource.id": f"prefect.work-queue.{work_queue.id}", 

687 "prefect.resource.name": work_queue.name, 

688 "prefect.resource.role": "work-queue", 

689 }, 

690 related=await _work_queue_work_pool_related(session, work_queue), 

691 id=uuid7(), 

692 ) 

693 

694 

695async def work_queue_deleted_event( 

696 session: AsyncSession, 

697 work_queue: "ORMWorkQueue", 

698 occurred: DateTime, 

699) -> Event: 

700 """Create an event for work queue deletion.""" 

701 return Event( 

702 occurred=occurred, 

703 event="prefect.work-queue.deleted", 

704 resource={ 

705 "prefect.resource.id": f"prefect.work-queue.{work_queue.id}", 

706 "prefect.resource.name": work_queue.name, 

707 "prefect.resource.role": "work-queue", 

708 }, 

709 related=await _work_queue_work_pool_related(session, work_queue), 

710 id=uuid7(), 

711 ) 

712 

713 

714def _get_recent_preceding_work_pool_event_id( 

715 work_pool: Optional["ORMWorkPool"], 

716) -> Optional[UUID]: 

717 """ 

718 Returns the preceding event ID if the work pool transitioned status 

719 recently to help ensure correct event ordering. 

720 """ 

721 if not work_pool: 

722 return None 

723 

724 time_since_last_event = timedelta(hours=24) 

725 if work_pool.last_transitioned_status_at: 

726 time_since_last_event = now("UTC") - work_pool.last_transitioned_status_at 

727 

728 return ( 

729 work_pool.last_status_event_id 

730 if time_since_last_event < timedelta(minutes=10) 

731 else None 

732 )