Coverage for /usr/local/lib/python3.12/site-packages/prefect/server/events/schemas/events.py: 65%
188 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
1import copy
2import fnmatch
3from collections import defaultdict
4from typing import (
5 TYPE_CHECKING,
6 Any,
7 ClassVar,
8 Dict,
9 Iterable,
10 List,
11 Mapping,
12 Optional,
13 Sequence,
14 Tuple,
15 Union,
16)
17from uuid import UUID
19from pydantic import (
20 AfterValidator,
21 AnyHttpUrl,
22 ConfigDict,
23 Field,
24 RootModel,
25 field_validator,
26 model_validator,
27)
28from typing_extensions import Annotated, Self
30import prefect.types._datetime
31from prefect.logging import get_logger
32from prefect.server.events.schemas.labelling import Labelled
33from prefect.server.utilities.schemas import PrefectBaseModel
34from prefect.settings import (
35 PREFECT_EVENTS_MAXIMUM_LABELS_PER_RESOURCE,
36 PREFECT_EVENTS_MAXIMUM_RELATED_RESOURCES,
37)
38from prefect.utilities.urls import url_for
40if TYPE_CHECKING: 40 ↛ 41line 40 didn't jump to line 41 because the condition on line 40 was never true
41 import logging
43 logger: "logging.Logger" = get_logger(__name__)
46class Resource(Labelled):
47 """An observable business object of interest to the user"""
49 @model_validator(mode="after")
50 def enforce_maximum_labels(self) -> Self:
51 if len(self.root) > PREFECT_EVENTS_MAXIMUM_LABELS_PER_RESOURCE.value(): 51 ↛ 52line 51 didn't jump to line 52 because the condition on line 51 was never true
52 raise ValueError(
53 "The maximum number of labels per resource "
54 f"is {PREFECT_EVENTS_MAXIMUM_LABELS_PER_RESOURCE.value()}"
55 )
57 return self
59 @model_validator(mode="after")
60 def requires_resource_id(self) -> Self:
61 if "prefect.resource.id" not in self.root:
62 raise ValueError("Resources must include the prefect.resource.id label")
63 if not self.root["prefect.resource.id"]: 63 ↛ 64line 63 didn't jump to line 64 because the condition on line 63 was never true
64 raise ValueError("The prefect.resource.id label must be non-empty")
66 return self
68 @property
69 def id(self) -> str:
70 return self["prefect.resource.id"]
72 @property
73 def name(self) -> Optional[str]:
74 return self.get("prefect.resource.name")
76 def prefect_object_id(self, kind: str) -> UUID:
77 """Extracts the UUID from an event's resource ID if it's the expected kind
78 of prefect resource"""
79 prefix = f"{kind}." if not kind.endswith(".") else kind
81 if not self.id.startswith(prefix):
82 raise ValueError(f"Resource ID {self.id} does not start with {prefix}")
84 return UUID(self.id[len(prefix) :])
87class RelatedResource(Resource):
88 """A Resource with a specific role in an Event"""
90 @model_validator(mode="after")
91 def requires_resource_role(self) -> Self:
92 if "prefect.resource.role" not in self.root: 92 ↛ 93line 92 didn't jump to line 93 because the condition on line 92 was never true
93 raise ValueError(
94 "Related Resources must include the prefect.resource.role label"
95 )
96 if not self.root["prefect.resource.role"]: 96 ↛ 97line 96 didn't jump to line 97 because the condition on line 96 was never true
97 raise ValueError("The prefect.resource.role label must be non-empty")
99 return self
101 @property
102 def role(self) -> str:
103 return self["prefect.resource.role"]
106def _validate_event_name_length(value: str) -> str:
107 from prefect.settings import PREFECT_SERVER_EVENTS_MAXIMUM_EVENT_NAME_LENGTH
109 if len(value) > PREFECT_SERVER_EVENTS_MAXIMUM_EVENT_NAME_LENGTH.value(): 109 ↛ 110line 109 didn't jump to line 110 because the condition on line 109 was never true
110 raise ValueError(
111 f"Event name must be at most {PREFECT_SERVER_EVENTS_MAXIMUM_EVENT_NAME_LENGTH.value()} characters"
112 )
113 return value
116class Event(PrefectBaseModel):
117 """The client-side view of an event that has happened to a Resource"""
119 occurred: prefect.types._datetime.DateTime = Field(
120 description="When the event happened from the sender's perspective",
121 )
122 event: Annotated[str, AfterValidator(_validate_event_name_length)] = Field(
123 description="The name of the event that happened",
124 )
125 resource: Resource = Field(
126 description="The primary Resource this event concerns",
127 )
128 related: list[RelatedResource] = Field(
129 default_factory=list,
130 description="A list of additional Resources involved in this event",
131 )
132 payload: dict[str, Any] = Field(
133 default_factory=dict,
134 description="An open-ended set of data describing what happened",
135 )
136 id: UUID = Field(
137 description="The client-provided identifier of this event",
138 )
139 follows: Optional[UUID] = Field(
140 default=None,
141 description=(
142 "The ID of an event that is known to have occurred prior to this one. "
143 "If set, this may be used to establish a more precise ordering of causally-"
144 "related events when they occur close enough together in time that the "
145 "system may receive them out-of-order."
146 ),
147 )
149 @property
150 def size_bytes(self) -> int:
151 return len(self.model_dump_json().encode())
153 @property
154 def involved_resources(self) -> Sequence[Resource]:
155 return [self.resource] + list(self.related)
157 @property
158 def resource_in_role(self) -> Mapping[str, RelatedResource]:
159 """Returns a mapping of roles to the first related resource in that role"""
160 return {related.role: related for related in reversed(self.related)}
162 @property
163 def resources_in_role(self) -> Mapping[str, Sequence[RelatedResource]]:
164 """Returns a mapping of roles to related resources in that role"""
165 resources: Dict[str, List[RelatedResource]] = defaultdict(list)
166 for related in self.related:
167 resources[related.role].append(related)
168 return resources
170 @field_validator("related")
171 @classmethod
172 def enforce_maximum_related_resources(
173 cls, value: List[RelatedResource]
174 ) -> List[RelatedResource]:
175 if len(value) > PREFECT_EVENTS_MAXIMUM_RELATED_RESOURCES.value(): 175 ↛ 176line 175 didn't jump to line 176 because the condition on line 175 was never true
176 raise ValueError(
177 "The maximum number of related resources "
178 f"is {PREFECT_EVENTS_MAXIMUM_RELATED_RESOURCES.value()}"
179 )
181 return value
183 def receive(
184 self, received: Optional[prefect.types._datetime.DateTime] = None
185 ) -> "ReceivedEvent":
186 kwargs = self.model_dump()
187 if received is not None: 187 ↛ 188line 187 didn't jump to line 188 because the condition on line 187 was never true
188 kwargs["received"] = received
189 return ReceivedEvent(**kwargs)
191 def find_resource_label(self, label: str) -> Optional[str]:
192 """Finds the value of the given label in this event's resource or one of its
193 related resources. If the label starts with `related:<role>:`, search for the
194 first matching label in a related resource with that role."""
195 directive, _, related_label = label.rpartition(":")
196 directive, _, role = directive.partition(":")
197 if directive == "related": 197 ↛ 198line 197 didn't jump to line 198 because the condition on line 197 was never true
198 for related in self.related:
199 if related.role == role:
200 return related.get(related_label)
201 return self.resource.get(label)
204class ReceivedEvent(Event):
205 """The server-side view of an event that has happened to a Resource after it has
206 been received by the server"""
208 model_config: ClassVar[ConfigDict] = ConfigDict(
209 extra="ignore", from_attributes=True
210 )
212 received: prefect.types._datetime.DateTime = Field(
213 default_factory=lambda: prefect.types._datetime.now("UTC"),
214 description="When the event was received by Prefect Cloud",
215 )
217 @property
218 def url(self) -> Optional[str]:
219 """Returns the UI URL for this event, allowing users to link to events
220 in automation templates without parsing date strings."""
221 return url_for(self, url_type="ui")
223 def as_database_row(self) -> dict[str, Any]:
224 row = self.model_dump()
225 row["resource_id"] = self.resource.id
226 row["recorded"] = prefect.types._datetime.now("UTC")
227 row["related_resource_ids"] = [related.id for related in self.related]
228 return row
230 def as_database_resource_rows(self) -> List[Dict[str, Any]]:
231 def without_id_and_role(resource: Resource) -> Dict[str, str]:
232 d: Dict[str, str] = resource.root.copy()
233 d.pop("prefect.resource.id", None)
234 d.pop("prefect.resource.role", None)
235 return d
237 return [
238 {
239 "occurred": self.occurred,
240 "resource_id": resource.id,
241 "resource_role": (
242 resource.role if isinstance(resource, RelatedResource) else ""
243 ),
244 "resource": without_id_and_role(resource),
245 "event_id": self.id,
246 }
247 for resource in [self.resource, *self.related]
248 ]
251def matches(expected: str, value: Optional[str]) -> bool:
252 """Returns true if the given value matches the expected string.
254 Args:
255 expected: A glob pattern to match against;
256 if it starts with an `!`, the pattern is negated.
257 value: The value of the label.
258 """
259 if value is None:
260 return False
262 is_positive = not expected.startswith("!")
263 expected = expected.removeprefix("!")
265 match = fnmatch.fnmatchcase(value, expected)
266 return match if is_positive else not match
269class ResourceSpecification(RootModel[Dict[str, Union[str, List[str]]]]):
270 def matches_every_resource(self) -> bool:
271 return len(self.root) == 0
273 def matches_every_resource_of_kind(self, prefix: str) -> bool:
274 if self.matches_every_resource():
275 return True
276 if len(self.root) == 1:
277 resource_id = self.root.get("prefect.resource.id")
278 if resource_id:
279 values = [resource_id] if isinstance(resource_id, str) else resource_id
280 return any(value == f"{prefix}.*" for value in values)
281 return False
283 def includes(self, candidates: Iterable[Resource]) -> bool:
284 if self.matches_every_resource():
285 return True
286 for candidate in candidates:
287 if self.matches(candidate): 287 ↛ 288line 287 didn't jump to line 288 because the condition on line 287 was never true
288 return True
289 return False
291 def matches(self, resource: Resource) -> bool:
292 for label, expected in self.items(): 292 ↛ 296line 292 didn't jump to line 296 because the loop on line 292 didn't complete
293 value = resource.get(label)
294 if not any(matches(candidate, value) for candidate in expected): 294 ↛ 292line 294 didn't jump to line 292 because the condition on line 294 was always true
295 return False
296 return True
298 def items(self) -> Iterable[Tuple[str, List[str]]]:
299 return [
300 (label, [value] if isinstance(value, str) else value)
301 for label, value in self.root.items()
302 ]
304 def __contains__(self, key: str) -> bool:
305 return key in self.root
307 def __getitem__(self, key: str) -> List[str]:
308 value = self.root[key]
309 if not value:
310 return []
311 if not isinstance(value, list):
312 value = [value]
313 return value
315 def pop(
316 self, key: str, default: Optional[Union[str, List[str]]] = None
317 ) -> Optional[List[str]]:
318 value = self.root.pop(key, default)
319 if not value: 319 ↛ 321line 319 didn't jump to line 321 because the condition on line 319 was always true
320 return []
321 if not isinstance(value, list):
322 value = [value]
323 return value
325 def get(
326 self, key: str, default: Optional[Union[str, List[str]]] = None
327 ) -> Optional[List[str]]:
328 value = self.root.get(key, default)
329 if not value:
330 return []
331 if not isinstance(value, list):
332 value = [value]
333 return value
335 def __len__(self) -> int:
336 return len(self.root)
338 def deepcopy(self) -> "ResourceSpecification":
339 return ResourceSpecification(root=copy.deepcopy(self.root))
342class EventPage(PrefectBaseModel):
343 """A single page of events returned from the API, with an optional link to the
344 next page of results"""
346 events: List[ReceivedEvent] = Field(
347 ..., description="The Events matching the query"
348 )
349 total: int = Field(..., description="The total number of matching Events")
350 next_page: Optional[AnyHttpUrl] = Field(
351 ..., description="The URL for the next page of results, if there are more"
352 )
355class EventCount(PrefectBaseModel):
356 """The count of events with the given filter value"""
358 value: str = Field(..., description="The value to use for filtering")
359 label: str = Field(..., description="The value to display for this count")
360 count: int = Field(..., description="The count of matching events")
361 start_time: prefect.types._datetime.DateTime = Field(
362 ..., description="The start time of this group of events"
363 )
364 end_time: prefect.types._datetime.DateTime = Field(
365 ..., description="The end time of this group of events"
366 )