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

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 

18 

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 

29 

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 

39 

40if TYPE_CHECKING: 40 ↛ 41line 40 didn't jump to line 41 because the condition on line 40 was never true

41 import logging 

42 

43 logger: "logging.Logger" = get_logger(__name__) 

44 

45 

46class Resource(Labelled): 

47 """An observable business object of interest to the user""" 

48 

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 ) 

56 

57 return self 

58 

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

65 

66 return self 

67 

68 @property 

69 def id(self) -> str: 

70 return self["prefect.resource.id"] 

71 

72 @property 

73 def name(self) -> Optional[str]: 

74 return self.get("prefect.resource.name") 

75 

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 

80 

81 if not self.id.startswith(prefix): 

82 raise ValueError(f"Resource ID {self.id} does not start with {prefix}") 

83 

84 return UUID(self.id[len(prefix) :]) 

85 

86 

87class RelatedResource(Resource): 

88 """A Resource with a specific role in an Event""" 

89 

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

98 

99 return self 

100 

101 @property 

102 def role(self) -> str: 

103 return self["prefect.resource.role"] 

104 

105 

106def _validate_event_name_length(value: str) -> str: 

107 from prefect.settings import PREFECT_SERVER_EVENTS_MAXIMUM_EVENT_NAME_LENGTH 

108 

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 

114 

115 

116class Event(PrefectBaseModel): 

117 """The client-side view of an event that has happened to a Resource""" 

118 

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 ) 

148 

149 @property 

150 def size_bytes(self) -> int: 

151 return len(self.model_dump_json().encode()) 

152 

153 @property 

154 def involved_resources(self) -> Sequence[Resource]: 

155 return [self.resource] + list(self.related) 

156 

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

161 

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 

169 

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 ) 

180 

181 return value 

182 

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) 

190 

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) 

202 

203 

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

207 

208 model_config: ClassVar[ConfigDict] = ConfigDict( 

209 extra="ignore", from_attributes=True 

210 ) 

211 

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 ) 

216 

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

222 

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 

229 

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 

236 

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 ] 

249 

250 

251def matches(expected: str, value: Optional[str]) -> bool: 

252 """Returns true if the given value matches the expected string. 

253 

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 

261 

262 is_positive = not expected.startswith("!") 

263 expected = expected.removeprefix("!") 

264 

265 match = fnmatch.fnmatchcase(value, expected) 

266 return match if is_positive else not match 

267 

268 

269class ResourceSpecification(RootModel[Dict[str, Union[str, List[str]]]]): 

270 def matches_every_resource(self) -> bool: 

271 return len(self.root) == 0 

272 

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 

282 

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 

290 

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 

297 

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 ] 

303 

304 def __contains__(self, key: str) -> bool: 

305 return key in self.root 

306 

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 

314 

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 

324 

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 

334 

335 def __len__(self) -> int: 

336 return len(self.root) 

337 

338 def deepcopy(self) -> "ResourceSpecification": 

339 return ResourceSpecification(root=copy.deepcopy(self.root)) 

340 

341 

342class EventPage(PrefectBaseModel): 

343 """A single page of events returned from the API, with an optional link to the 

344 next page of results""" 

345 

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 ) 

353 

354 

355class EventCount(PrefectBaseModel): 

356 """The count of events with the given filter value""" 

357 

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 )