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

1"""Typed payloads and builders for `prefect.<object>.{created,updated,deleted}` 

2lifecycle events. 

3 

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. 

8 

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

13 

14from typing import TYPE_CHECKING, Any, ClassVar, Dict, List, Optional 

15from uuid import UUID 

16 

17from pydantic import ConfigDict, Field 

18 

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 

24 

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 

36 

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

38 

39 

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 ) 

60 

61 

62class VariableEventPayload(PrefectBaseModel): 

63 """The payload of a variable lifecycle event: its name, value, and tags.""" 

64 

65 model_config: ClassVar[ConfigDict] = ConfigDict(from_attributes=True) 

66 

67 name: str 

68 value: StrictVariableValue 

69 tags: List[str] = Field(default_factory=list) 

70 

71 

72def variable_created_event(variable: "ORMVariable", occurred: DateTime) -> Event: 

73 """Create an event for variable creation.""" 

74 return _variable_event("created", variable, occurred) 

75 

76 

77def variable_updated_event(variable: "ORMVariable", occurred: DateTime) -> Event: 

78 """Create an event for variable updates.""" 

79 return _variable_event("updated", variable, occurred) 

80 

81 

82def variable_deleted_event(variable: "ORMVariable", occurred: DateTime) -> Event: 

83 """Create an event for variable deletion.""" 

84 return _variable_event("deleted", variable, occurred) 

85 

86 

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 ) 

96 

97 

98class FlowEventPayload(PrefectBaseModel): 

99 """The payload of a flow lifecycle event: its name, tags, and labels.""" 

100 

101 model_config: ClassVar[ConfigDict] = ConfigDict(from_attributes=True) 

102 

103 name: str 

104 tags: List[str] = Field(default_factory=list) 

105 labels: Optional[KeyValueLabels] = Field(default_factory=dict) 

106 

107 

108def flow_created_event(flow: "ORMFlow", occurred: DateTime) -> Event: 

109 """Create an event for flow creation.""" 

110 return _flow_event("created", flow, occurred) 

111 

112 

113def flow_updated_event(flow: "ORMFlow", occurred: DateTime) -> Event: 

114 """Create an event for flow updates.""" 

115 return _flow_event("updated", flow, occurred) 

116 

117 

118def flow_deleted_event(flow: "ORMFlow", occurred: DateTime) -> Event: 

119 """Create an event for flow deletion.""" 

120 return _flow_event("deleted", flow, occurred) 

121 

122 

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 ) 

132 

133 

134class BlockTypeEventPayload(PrefectBaseModel): 

135 """The payload of a block type lifecycle event: its identity, presentation, 

136 and protection.""" 

137 

138 model_config: ClassVar[ConfigDict] = ConfigDict(from_attributes=True) 

139 

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 

147 

148 

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) 

152 

153 

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) 

157 

158 

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) 

162 

163 

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 ) 

177 

178 

179class BlockDocumentBlockType(PrefectBaseModel): 

180 """The block type a block document belongs to, nested on its block schema.""" 

181 

182 id: UUID 

183 name: str 

184 

185 

186class BlockSchemaEventPayload(PrefectBaseModel): 

187 """The schema a block document conforms to, with its block type nested.""" 

188 

189 capabilities: List[str] = Field(default_factory=list) 

190 block_type: Optional[BlockDocumentBlockType] = None 

191 

192 

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

197 

198 name: Optional[str] = None 

199 data: Dict[str, str] = Field(default_factory=dict) 

200 block_schema: Optional[BlockSchemaEventPayload] = None 

201 

202 

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) 

208 

209 

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) 

215 

216 

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) 

222 

223 

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 ) 

231 

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 ) 

241 

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 ) 

251 

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 

266 

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 ) 

276 

277 

278class ConcurrencyLimitV2EventPayload(PrefectBaseModel): 

279 """The payload of a global concurrency limit lifecycle event: its name, 

280 limit, active flag, and slot decay rate.""" 

281 

282 model_config: ClassVar[ConfigDict] = ConfigDict(from_attributes=True) 

283 

284 name: str 

285 limit: int 

286 active: bool 

287 slot_decay_per_second: float 

288 

289 

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) 

295 

296 

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) 

302 

303 

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) 

309 

310 

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 ) 

324 

325 

326class ConcurrencyLimitEventPayload(PrefectBaseModel): 

327 """The payload of a tag-based (v1) concurrency limit lifecycle event: the tag 

328 it applies to and its limit.""" 

329 

330 model_config: ClassVar[ConfigDict] = ConfigDict(from_attributes=True) 

331 

332 tag: str 

333 concurrency_limit: int 

334 

335 

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) 

341 

342 

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) 

348 

349 

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) 

355 

356 

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 ) 

370 

371 

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

376 

377 model_config: ClassVar[ConfigDict] = ConfigDict(from_attributes=True) 

378 

379 key: str 

380 type: Optional[str] = None 

381 data: Optional[Any] = None 

382 description: Optional[str] = None 

383 

384 

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) 

390 

391 

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) 

397 

398 

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) 

404 

405 

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 ) 

435 

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 ) 

447 

448 

449def automation_created_event(automation: "Automation", occurred: DateTime) -> Event: 

450 """Create an event for automation creation.""" 

451 return _automation_event("created", automation, occurred) 

452 

453 

454def automation_updated_event(automation: "Automation", occurred: DateTime) -> Event: 

455 """Create an event for automation updates.""" 

456 return _automation_event("updated", automation, occurred) 

457 

458 

459def automation_deleted_event(automation: "Automation", occurred: DateTime) -> Event: 

460 """Create an event for automation deletion.""" 

461 return _automation_event("deleted", automation, occurred) 

462 

463 

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 )