Coverage for /usr/local/lib/python3.12/site-packages/prefect/server/events/schemas/lifecycle.py: 92%
141 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"""Typed payloads and builders for `prefect.<object>.{created,updated,deleted}`
2lifecycle events.
4Each payload is the shape of the event's `payload` field for a domain object —
5what a consumer sees on the event without reading the object back from the API.
6The shapes mirror Prefect Cloud's lifecycle-event payloads so the two emit the
7same thing.
9The builders here are pure functions of an ORM object and a timestamp, with no
10dependency on the model layer, so any model module can import them without
11risking a circular import.
12"""
14from typing import TYPE_CHECKING, Any, ClassVar, Dict, List, Optional
15from uuid import UUID
17from pydantic import ConfigDict, Field
19from prefect._internal.uuid7 import uuid7
20from prefect.server.events.schemas.events import Event
21from prefect.server.utilities.schemas.bases import PrefectBaseModel
22from prefect.types import KeyValueLabels, StrictVariableValue
23from prefect.types._datetime import DateTime
25if TYPE_CHECKING: 25 ↛ 26line 25 didn't jump to line 26 because the condition on line 25 was never true
26 from prefect.server.database.orm_models import (
27 BlockType,
28 ConcurrencyLimit,
29 ConcurrencyLimitV2,
30 ORMArtifactCollection,
31 ORMFlow,
32 ORMVariable,
33 )
34 from prefect.server.events.schemas.automations import Automation
35 from prefect.server.schemas.core import BlockDocument
37RelatedResourceList = List[Dict[str, str]]
40def _lifecycle_event(
41 kind: str,
42 action: str,
43 resource_id: str,
44 resource_name: Optional[str],
45 payload: Dict[str, Any],
46 occurred: DateTime,
47 related: Optional[RelatedResourceList] = None,
48) -> Event:
49 resource: Dict[str, str] = {"prefect.resource.id": resource_id}
50 if resource_name is not None: 50 ↛ 52line 50 didn't jump to line 52 because the condition on line 50 was always true
51 resource["prefect.resource.name"] = resource_name
52 return Event(
53 occurred=occurred,
54 event=f"prefect.{kind}.{action}",
55 resource=resource,
56 related=related or [],
57 payload=payload,
58 id=uuid7(),
59 )
62class VariableEventPayload(PrefectBaseModel):
63 """The payload of a variable lifecycle event: its name, value, and tags."""
65 model_config: ClassVar[ConfigDict] = ConfigDict(from_attributes=True)
67 name: str
68 value: StrictVariableValue
69 tags: List[str] = Field(default_factory=list)
72def variable_created_event(variable: "ORMVariable", occurred: DateTime) -> Event:
73 """Create an event for variable creation."""
74 return _variable_event("created", variable, occurred)
77def variable_updated_event(variable: "ORMVariable", occurred: DateTime) -> Event:
78 """Create an event for variable updates."""
79 return _variable_event("updated", variable, occurred)
82def variable_deleted_event(variable: "ORMVariable", occurred: DateTime) -> Event:
83 """Create an event for variable deletion."""
84 return _variable_event("deleted", variable, occurred)
87def _variable_event(action: str, variable: "ORMVariable", occurred: DateTime) -> Event:
88 return _lifecycle_event(
89 kind="variable",
90 action=action,
91 resource_id=f"prefect.variable.{variable.id}",
92 resource_name=variable.name,
93 payload=VariableEventPayload.model_validate(variable).model_dump(mode="json"),
94 occurred=occurred,
95 )
98class FlowEventPayload(PrefectBaseModel):
99 """The payload of a flow lifecycle event: its name, tags, and labels."""
101 model_config: ClassVar[ConfigDict] = ConfigDict(from_attributes=True)
103 name: str
104 tags: List[str] = Field(default_factory=list)
105 labels: Optional[KeyValueLabels] = Field(default_factory=dict)
108def flow_created_event(flow: "ORMFlow", occurred: DateTime) -> Event:
109 """Create an event for flow creation."""
110 return _flow_event("created", flow, occurred)
113def flow_updated_event(flow: "ORMFlow", occurred: DateTime) -> Event:
114 """Create an event for flow updates."""
115 return _flow_event("updated", flow, occurred)
118def flow_deleted_event(flow: "ORMFlow", occurred: DateTime) -> Event:
119 """Create an event for flow deletion."""
120 return _flow_event("deleted", flow, occurred)
123def _flow_event(action: str, flow: "ORMFlow", occurred: DateTime) -> Event:
124 return _lifecycle_event(
125 kind="flow",
126 action=action,
127 resource_id=f"prefect.flow.{flow.id}",
128 resource_name=flow.name,
129 payload=FlowEventPayload.model_validate(flow).model_dump(mode="json"),
130 occurred=occurred,
131 )
134class BlockTypeEventPayload(PrefectBaseModel):
135 """The payload of a block type lifecycle event: its identity, presentation,
136 and protection."""
138 model_config: ClassVar[ConfigDict] = ConfigDict(from_attributes=True)
140 name: str
141 slug: str
142 logo_url: Optional[str] = None
143 documentation_url: Optional[str] = None
144 description: Optional[str] = None
145 code_example: Optional[str] = None
146 is_protected: bool = False
149def block_type_created_event(block_type: "BlockType", occurred: DateTime) -> Event:
150 """Create an event for block type creation."""
151 return _block_type_event("created", block_type, occurred)
154def block_type_updated_event(block_type: "BlockType", occurred: DateTime) -> Event:
155 """Create an event for block type updates."""
156 return _block_type_event("updated", block_type, occurred)
159def block_type_deleted_event(block_type: "BlockType", occurred: DateTime) -> Event:
160 """Create an event for block type deletion."""
161 return _block_type_event("deleted", block_type, occurred)
164def _block_type_event(
165 action: str, block_type: "BlockType", occurred: DateTime
166) -> Event:
167 return _lifecycle_event(
168 kind="block-type",
169 action=action,
170 resource_id=f"prefect.block-type.{block_type.id}",
171 resource_name=block_type.name,
172 payload=BlockTypeEventPayload.model_validate(block_type).model_dump(
173 mode="json"
174 ),
175 occurred=occurred,
176 )
179class BlockDocumentBlockType(PrefectBaseModel):
180 """The block type a block document belongs to, nested on its block schema."""
182 id: UUID
183 name: str
186class BlockSchemaEventPayload(PrefectBaseModel):
187 """The schema a block document conforms to, with its block type nested."""
189 capabilities: List[str] = Field(default_factory=list)
190 block_type: Optional[BlockDocumentBlockType] = None
193class BlockDocumentEventPayload(PrefectBaseModel):
194 """The payload of a block document lifecycle event: its name, the non-secret
195 string values of its data, and the nested block schema (capabilities and
196 block type)."""
198 name: Optional[str] = None
199 data: Dict[str, str] = Field(default_factory=dict)
200 block_schema: Optional[BlockSchemaEventPayload] = None
203def block_document_created_event(
204 block_document: "BlockDocument", occurred: DateTime
205) -> Event:
206 """Create an event for block document creation."""
207 return _block_document_event("created", block_document, occurred)
210def block_document_updated_event(
211 block_document: "BlockDocument", occurred: DateTime
212) -> Event:
213 """Create an event for block document updates."""
214 return _block_document_event("updated", block_document, occurred)
217def block_document_deleted_event(
218 block_document: "BlockDocument", occurred: DateTime
219) -> Event:
220 """Create an event for block document deletion."""
221 return _block_document_event("deleted", block_document, occurred)
224def _block_document_event(
225 action: str, block_document: "BlockDocument", occurred: DateTime
226) -> Event:
227 schema = block_document.block_schema
228 secret_keys = set(
229 (schema.fields.get("secret_fields") or []) if schema is not None else []
230 )
232 block_schema_payload: Optional[BlockSchemaEventPayload] = None
233 if schema is not None: 233 ↛ 242line 233 didn't jump to line 242 because the condition on line 233 was always true
234 block_schema_payload = BlockSchemaEventPayload(
235 capabilities=schema.capabilities,
236 block_type=BlockDocumentBlockType(
237 id=block_document.block_type_id,
238 name=block_document.block_type_name or "",
239 ),
240 )
242 payload = BlockDocumentEventPayload(
243 name=block_document.name,
244 data={
245 key: value
246 for key, value in block_document.data.items()
247 if key not in secret_keys and isinstance(value, str)
248 },
249 block_schema=block_schema_payload,
250 )
252 related: RelatedResourceList = [
253 {
254 "prefect.resource.id": f"prefect.block-type.{block_document.block_type_id}",
255 "prefect.resource.role": "block-type",
256 },
257 {
258 "prefect.resource.id": (
259 f"prefect.block-schema.{block_document.block_schema_id}"
260 ),
261 "prefect.resource.role": "block-schema",
262 },
263 ]
264 if block_document.block_type_name: 264 ↛ 267line 264 didn't jump to line 267 because the condition on line 264 was always true
265 related[0]["prefect.resource.name"] = block_document.block_type_name
267 return _lifecycle_event(
268 kind="block-document",
269 action=action,
270 resource_id=f"prefect.block-document.{block_document.id}",
271 resource_name=block_document.name,
272 payload=payload.model_dump(mode="json"),
273 occurred=occurred,
274 related=related,
275 )
278class ConcurrencyLimitV2EventPayload(PrefectBaseModel):
279 """The payload of a global concurrency limit lifecycle event: its name,
280 limit, active flag, and slot decay rate."""
282 model_config: ClassVar[ConfigDict] = ConfigDict(from_attributes=True)
284 name: str
285 limit: int
286 active: bool
287 slot_decay_per_second: float
290def concurrency_limit_v2_created_event(
291 concurrency_limit: "ConcurrencyLimitV2", occurred: DateTime
292) -> Event:
293 """Create an event for global concurrency limit creation."""
294 return _concurrency_limit_v2_event("created", concurrency_limit, occurred)
297def concurrency_limit_v2_updated_event(
298 concurrency_limit: "ConcurrencyLimitV2", occurred: DateTime
299) -> Event:
300 """Create an event for global concurrency limit updates."""
301 return _concurrency_limit_v2_event("updated", concurrency_limit, occurred)
304def concurrency_limit_v2_deleted_event(
305 concurrency_limit: "ConcurrencyLimitV2", occurred: DateTime
306) -> Event:
307 """Create an event for global concurrency limit deletion."""
308 return _concurrency_limit_v2_event("deleted", concurrency_limit, occurred)
311def _concurrency_limit_v2_event(
312 action: str, concurrency_limit: "ConcurrencyLimitV2", occurred: DateTime
313) -> Event:
314 return _lifecycle_event(
315 kind="concurrency-limit",
316 action=action,
317 resource_id=f"prefect.concurrency-limit.{concurrency_limit.id}",
318 resource_name=concurrency_limit.name,
319 payload=ConcurrencyLimitV2EventPayload.model_validate(
320 concurrency_limit
321 ).model_dump(mode="json"),
322 occurred=occurred,
323 )
326class ConcurrencyLimitEventPayload(PrefectBaseModel):
327 """The payload of a tag-based (v1) concurrency limit lifecycle event: the tag
328 it applies to and its limit."""
330 model_config: ClassVar[ConfigDict] = ConfigDict(from_attributes=True)
332 tag: str
333 concurrency_limit: int
336def concurrency_limit_created_event(
337 concurrency_limit: "ConcurrencyLimit", occurred: DateTime
338) -> Event:
339 """Create an event for tag-based concurrency limit creation."""
340 return _concurrency_limit_event("created", concurrency_limit, occurred)
343def concurrency_limit_updated_event(
344 concurrency_limit: "ConcurrencyLimit", occurred: DateTime
345) -> Event:
346 """Create an event for tag-based concurrency limit updates."""
347 return _concurrency_limit_event("updated", concurrency_limit, occurred)
350def concurrency_limit_deleted_event(
351 concurrency_limit: "ConcurrencyLimit", occurred: DateTime
352) -> Event:
353 """Create an event for tag-based concurrency limit deletion."""
354 return _concurrency_limit_event("deleted", concurrency_limit, occurred)
357def _concurrency_limit_event(
358 action: str, concurrency_limit: "ConcurrencyLimit", occurred: DateTime
359) -> Event:
360 return _lifecycle_event(
361 kind="concurrency-limit",
362 action=action,
363 resource_id=f"prefect.concurrency-limit.{concurrency_limit.id}",
364 resource_name=concurrency_limit.tag,
365 payload=ConcurrencyLimitEventPayload.model_validate(
366 concurrency_limit
367 ).model_dump(mode="json"),
368 occurred=occurred,
369 )
372class ArtifactCollectionEventPayload(PrefectBaseModel):
373 """The payload of an artifact collection lifecycle event: the collection key,
374 the artifact type, and its data and description. The latest artifact and its
375 flow/task runs ride the event's `related` resources, not the payload."""
377 model_config: ClassVar[ConfigDict] = ConfigDict(from_attributes=True)
379 key: str
380 type: Optional[str] = None
381 data: Optional[Any] = None
382 description: Optional[str] = None
385def artifact_collection_created_event(
386 artifact_collection: "ORMArtifactCollection", occurred: DateTime
387) -> Event:
388 """Create an event for artifact collection creation."""
389 return _artifact_collection_event("created", artifact_collection, occurred)
392def artifact_collection_updated_event(
393 artifact_collection: "ORMArtifactCollection", occurred: DateTime
394) -> Event:
395 """Create an event for artifact collection updates."""
396 return _artifact_collection_event("updated", artifact_collection, occurred)
399def artifact_collection_deleted_event(
400 artifact_collection: "ORMArtifactCollection", occurred: DateTime
401) -> Event:
402 """Create an event for artifact collection deletion."""
403 return _artifact_collection_event("deleted", artifact_collection, occurred)
406def _artifact_collection_event(
407 action: str, artifact_collection: "ORMArtifactCollection", occurred: DateTime
408) -> Event:
409 related: RelatedResourceList = [
410 {
411 "prefect.resource.id": (
412 f"prefect.artifact.{artifact_collection.latest_id}"
413 ),
414 "prefect.resource.role": "latest",
415 }
416 ]
417 if artifact_collection.flow_run_id:
418 related.append(
419 {
420 "prefect.resource.id": (
421 f"prefect.flow-run.{artifact_collection.flow_run_id}"
422 ),
423 "prefect.resource.role": "flow-run",
424 }
425 )
426 if artifact_collection.task_run_id:
427 related.append(
428 {
429 "prefect.resource.id": (
430 f"prefect.task-run.{artifact_collection.task_run_id}"
431 ),
432 "prefect.resource.role": "task-run",
433 }
434 )
436 return _lifecycle_event(
437 kind="artifact-collection",
438 action=action,
439 resource_id=f"prefect.artifact-collection.{artifact_collection.id}",
440 resource_name=artifact_collection.key,
441 payload=ArtifactCollectionEventPayload.model_validate(
442 artifact_collection
443 ).model_dump(mode="json"),
444 occurred=occurred,
445 related=related,
446 )
449def automation_created_event(automation: "Automation", occurred: DateTime) -> Event:
450 """Create an event for automation creation."""
451 return _automation_event("created", automation, occurred)
454def automation_updated_event(automation: "Automation", occurred: DateTime) -> Event:
455 """Create an event for automation updates."""
456 return _automation_event("updated", automation, occurred)
459def automation_deleted_event(automation: "Automation", occurred: DateTime) -> Event:
460 """Create an event for automation deletion."""
461 return _automation_event("deleted", automation, occurred)
464def _automation_event(
465 action: str, automation: "Automation", occurred: DateTime
466) -> Event:
467 return _lifecycle_event(
468 kind="automation",
469 action=action,
470 resource_id=f"prefect.automation.{automation.id}",
471 resource_name=automation.name,
472 payload=automation.model_dump(mode="json"),
473 occurred=occurred,
474 )