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
« 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
5from cachetools import TTLCache
6from sqlalchemy.ext.asyncio import AsyncSession
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
27ResourceData = Dict[str, Dict[str, Any]]
28RelatedResourceList = List[Dict[str, str]]
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
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)
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 )
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 )
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 )
105 work_pool = work_queue.work_pool if work_queue is not None else None
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 )
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
119 return _resource_data_as_related_resources(
120 resource_data,
121 excluded_kinds=["flow-run"],
122 ) + _provenance_as_related_resources(flow_run.created_by)
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 }
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()
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 = []
204 for kind, data in resource_data.items():
205 tags |= set(data.get("tags", []))
207 if kind in excluded_kinds or not data:
208 continue
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 }
216 if kind == "work-pool":
217 related_resource["prefect.work-pool.type"] = data["type"]
219 related.append(related_resource)
221 related += [
222 {
223 "prefect.resource.id": f"prefect.tag.{tag}",
224 "prefect.resource.role": "tag",
225 }
226 for tag in sorted(tags)
227 ]
229 return related
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 []
238 resource_id: str
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 []
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]
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
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
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)
291 return False
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 = []
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 )
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 )
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 )
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 )
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 )
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 )
376 return related
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
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 )
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 )
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 )
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 )
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]] = []
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 )
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 )
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 )
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
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 )
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 )
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]] = []
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 )
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 )
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 = []
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 )
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 )
614 return related
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 }
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 )
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 )
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
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 )
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 )
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
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
728 return (
729 work_pool.last_status_event_id
730 if time_since_last_event < timedelta(minutes=10)
731 else None
732 )