Coverage for /usr/local/lib/python3.12/site-packages/prefect/server/events/schemas/automations.py: 80%

312 statements  

« prev     ^ index     » next       coverage.py v7.15.2, created at 2026-10-07 02:04 +0000

1from __future__ import annotations 

2 

3import abc 

4import re 

5import weakref 

6from datetime import timedelta 

7from typing import ( 

8 TYPE_CHECKING, 

9 Any, 

10 Dict, 

11 List, 

12 Literal, 

13 Optional, 

14 Sequence, 

15 Set, 

16 Tuple, 

17 Type, 

18 TypeVar, 

19 Union, 

20) 

21from uuid import UUID, uuid4 

22 

23from pydantic import ( 

24 Field, 

25 PrivateAttr, 

26 field_validator, 

27 model_validator, 

28) 

29from typing_extensions import Self, TypeAlias 

30 

31from prefect._internal.uuid7 import uuid7 

32from prefect.logging import get_logger 

33from prefect.server.events.actions import ServerActionTypes 

34from prefect.server.events.schemas.events import ( 

35 ReceivedEvent, 

36 RelatedResource, 

37 Resource, 

38 ResourceSpecification, 

39 matches, 

40) 

41from prefect.server.schemas.actions import ActionBaseModel 

42from prefect.server.utilities.schemas import ORMBaseModel, PrefectBaseModel 

43from prefect.types import DateTime 

44from prefect.utilities.collections import AutoEnum 

45 

46# Dict form accepted by ResourceSpecification via Pydantic coercion 

47ResourceSpecificationDict: TypeAlias = Dict[str, Union[str, List[str]]] 

48 

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

50 import logging 

51 

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

53 

54 

55class Posture(AutoEnum): 

56 Reactive = "Reactive" 

57 Proactive = "Proactive" 

58 Metric = "Metric" 

59 

60 

61class TriggerState(AutoEnum): 

62 Triggered = "Triggered" 

63 Resolved = "Resolved" 

64 

65 

66class Trigger(PrefectBaseModel, abc.ABC): 

67 """ 

68 Base class describing a set of criteria that must be satisfied in order to trigger 

69 an automation. 

70 """ 

71 

72 type: str 

73 

74 id: UUID = Field(default_factory=uuid4, description="The unique ID of this trigger") 

75 

76 _automation: Optional[weakref.ref[Any]] = PrivateAttr(None) 

77 _parent: Optional[weakref.ref[Any]] = PrivateAttr(None) 

78 

79 @property 

80 def automation(self) -> "Automation": 

81 assert self._automation is not None, "Trigger._automation has not been set" 

82 value = self._automation() 

83 assert value is not None, "Trigger._automation has been garbage collected" 

84 return value 

85 

86 @property 

87 def parent(self) -> "Union[Trigger, Automation]": 

88 assert self._parent is not None, "Trigger._parent has not been set" 

89 value = self._parent() 

90 assert value is not None, "Trigger._parent has been garbage collected" 

91 return value 

92 

93 def _set_parent(self, value: "Union[Trigger, Automation]"): 

94 if isinstance(value, Automation): 

95 self._automation = weakref.ref(value) 

96 self._parent = self._automation 

97 elif isinstance(value, Trigger): 97 ↛ 101line 97 didn't jump to line 101 because the condition on line 97 was always true

98 self._parent = weakref.ref(value) 

99 self._automation = value._automation 

100 else: # pragma: no cover 

101 raise ValueError("parent must be an Automation or a Trigger") 

102 

103 def reset_ids(self) -> None: 

104 """Resets the ID of this trigger and all of its children""" 

105 self.id = uuid4() 

106 for trigger in self.all_triggers(): 

107 trigger.id = uuid4() 

108 

109 def all_triggers(self) -> Sequence[Trigger]: 

110 """Returns all triggers within this trigger""" 

111 return [self] 

112 

113 @abc.abstractmethod 

114 def create_automation_state_change_event( 114 ↛ exitline 114 didn't return from function 'create_automation_state_change_event' because

115 self, firing: "Firing", trigger_state: TriggerState 

116 ) -> ReceivedEvent: ... 

117 

118 

119class CompositeTrigger(Trigger, abc.ABC): 

120 """ 

121 Requires some number of triggers to have fired within the given time period. 

122 """ 

123 

124 type: Literal["compound", "sequence"] 

125 triggers: List["ServerTriggerTypes"] 

126 within: Optional[timedelta] 

127 

128 def create_automation_state_change_event( 

129 self, firing: Firing, trigger_state: TriggerState 

130 ) -> ReceivedEvent: 

131 """Returns a ReceivedEvent for an automation state change 

132 into a triggered or resolved state.""" 

133 automation = firing.trigger.automation 

134 triggering_event = firing.triggering_event 

135 return ReceivedEvent( 

136 occurred=firing.triggered, 

137 event=f"prefect.automation.{trigger_state.value.lower()}", 

138 resource={ 

139 "prefect.resource.id": f"prefect.automation.{automation.id}", 

140 "prefect.resource.name": automation.name, 

141 }, 

142 related=( 

143 [ 

144 { 

145 "prefect.resource.id": f"prefect.event.{triggering_event.id}", 

146 "prefect.resource.role": "triggering-event", 

147 } 

148 ] 

149 if triggering_event 

150 else [] 

151 ), 

152 payload={ 

153 "triggering_labels": firing.triggering_labels, 

154 "triggering_event": ( 

155 triggering_event.model_dump(mode="json") 

156 if triggering_event 

157 else None 

158 ), 

159 }, 

160 id=uuid7(), 

161 ) 

162 

163 def _set_parent(self, value: "Union[Trigger , Automation]"): 

164 super()._set_parent(value) 

165 for trigger in self.triggers: 

166 trigger._set_parent(self) 

167 

168 def all_triggers(self) -> Sequence[Trigger]: 

169 return [self] + [t for child in self.triggers for t in child.all_triggers()] 

170 

171 @property 

172 def child_trigger_ids(self) -> List[UUID]: 

173 return [trigger.id for trigger in self.triggers] 

174 

175 @property 

176 def num_expected_firings(self) -> int: 

177 return len(self.triggers) 

178 

179 @abc.abstractmethod 

180 def ready_to_fire(self, firings: Sequence["Firing"]) -> bool: ... 180 ↛ exitline 180 didn't return from function 'ready_to_fire' because

181 

182 

183class CompoundTrigger(CompositeTrigger): 

184 """A composite trigger that requires some number of triggers to have 

185 fired within the given time period""" 

186 

187 type: Literal["compound"] = "compound" 

188 require: Union[int, Literal["any", "all"]] 

189 

190 @property 

191 def num_expected_firings(self) -> int: 

192 if self.require == "any": 

193 return 1 

194 elif self.require == "all": 

195 return len(self.triggers) 

196 else: 

197 return int(self.require) 

198 

199 def ready_to_fire(self, firings: Sequence["Firing"]) -> bool: 

200 return len(firings) >= self.num_expected_firings 

201 

202 @model_validator(mode="after") 

203 def validate_require(self) -> Self: 

204 if isinstance(self.require, int): 204 ↛ 212line 204 didn't jump to line 212 because the condition on line 204 was always true

205 if self.require < 1: 

206 raise ValueError("require must be at least 1") 

207 if self.require > len(self.triggers): 207 ↛ 212line 207 didn't jump to line 212 because the condition on line 207 was always true

208 raise ValueError( 

209 "require must be less than or equal to the number of triggers" 

210 ) 

211 

212 return self 

213 

214 

215class SequenceTrigger(CompositeTrigger): 

216 """A composite trigger that requires some number of triggers to have fired 

217 within the given time period in a specific order""" 

218 

219 type: Literal["sequence"] = "sequence" 

220 

221 @property 

222 def expected_firing_order(self) -> List[UUID]: 

223 return [trigger.id for trigger in self.triggers] 

224 

225 def ready_to_fire(self, firings: Sequence["Firing"]) -> bool: 

226 actual_firing_order = [ 

227 f.trigger.id for f in sorted(firings, key=lambda f: f.triggered) 

228 ] 

229 return actual_firing_order == self.expected_firing_order 

230 

231 

232class ResourceTrigger(Trigger, abc.ABC): 

233 """ 

234 Base class for triggers that may filter by the labels of resources. 

235 """ 

236 

237 type: str 

238 

239 match: Union[ResourceSpecification, ResourceSpecificationDict] = Field( 

240 default_factory=lambda: ResourceSpecification.model_validate({}), 

241 description="Labels for resources which this trigger will match.", 

242 ) 

243 match_related: Union[ 

244 ResourceSpecification, 

245 ResourceSpecificationDict, 

246 list[Union[ResourceSpecification, ResourceSpecificationDict]], 

247 ] = Field( 

248 default_factory=lambda: ResourceSpecification.model_validate({}), 

249 description="Labels for related resources which this trigger will match.", 

250 ) 

251 

252 @field_validator("match", mode="before") 

253 @classmethod 

254 def coerce_match(cls, v: Any) -> Any: 

255 if isinstance(v, dict): 

256 return ResourceSpecification.model_validate(v) 

257 return v 

258 

259 @field_validator("match_related", mode="before") 

260 @classmethod 

261 def coerce_match_related(cls, v: Any) -> Any: 

262 if isinstance(v, dict): 

263 return ResourceSpecification.model_validate(v) 

264 if isinstance(v, list): 

265 return [ 

266 ResourceSpecification.model_validate(item) 

267 if isinstance(item, dict) 

268 else item 

269 for item in v 

270 ] 

271 return v 

272 

273 def covers_resources( 

274 self, resource: Resource, related: Sequence[RelatedResource] 

275 ) -> bool: 

276 if not self.match.includes([resource]): 

277 return False 

278 

279 match_related = self.match_related 

280 if not isinstance(match_related, list): 

281 match_related = [match_related] 

282 

283 if not all(match.includes(related) for match in match_related): 

284 return False 

285 

286 return True 

287 

288 

289class EventTrigger(ResourceTrigger): 

290 """ 

291 A trigger that fires based on the presence or absence of events within a given 

292 period of time. 

293 """ 

294 

295 type: Literal["event"] = "event" 

296 

297 after: Set[str] = Field( 

298 default_factory=set, 

299 description=( 

300 "The event(s) which must first been seen to fire this trigger. If " 

301 "empty, then fire this trigger immediately. Events may include " 

302 "trailing wildcards, like `prefect.flow-run.*`" 

303 ), 

304 ) 

305 expect: Set[str] = Field( 

306 default_factory=set, 

307 description=( 

308 "The event(s) this trigger is expecting to see. If empty, this " 

309 "trigger will match any event. Events may include trailing wildcards, " 

310 "like `prefect.flow-run.*`" 

311 ), 

312 ) 

313 

314 for_each: Set[str] = Field( 

315 default_factory=set, 

316 description=( 

317 "Evaluate the trigger separately for each distinct value of these labels " 

318 "on the resource. By default, labels refer to the primary resource of the " 

319 "triggering event. You may also refer to labels from related " 

320 "resources by specifying `related:<role>:<label>`. This will use the " 

321 "value of that label for the first related resource in that role. For " 

322 'example, `"for_each": ["related:flow:prefect.resource.id"]` would ' 

323 "evaluate the trigger for each flow." 

324 ), 

325 ) 

326 posture: Literal[Posture.Reactive, Posture.Proactive] = Field( # type: ignore[valid-type] 

327 ..., 

328 description=( 

329 "The posture of this trigger, either Reactive or Proactive. Reactive " 

330 "triggers respond to the _presence_ of the expected events, while " 

331 "Proactive triggers respond to the _absence_ of those expected events." 

332 ), 

333 ) 

334 threshold: int = Field( 

335 1, 

336 description=( 

337 "The number of events required for this trigger to fire (for " 

338 "Reactive triggers), or the number of events expected (for Proactive " 

339 "triggers)" 

340 ), 

341 ) 

342 within: timedelta = Field( 

343 timedelta(seconds=0), 

344 ge=timedelta(seconds=0), 

345 description=( 

346 "The time period over which the events must occur. For Reactive triggers, " 

347 "this may be as low as 0 seconds, but must be at least 10 seconds for " 

348 "Proactive triggers" 

349 ), 

350 ) 

351 

352 @model_validator(mode="before") 

353 @classmethod 

354 def enforce_minimum_within_for_proactive_triggers( 

355 cls, data: Dict[str, Any] | Any 

356 ) -> Dict[str, Any]: 

357 if not isinstance(data, dict): 

358 return data 

359 

360 if "within" in data and data["within"] is None: 

361 raise ValueError("`within` should be a valid timedelta") 

362 

363 posture: Optional[Posture] = data.get("posture") 

364 within: Optional[timedelta] = data.get("within") 

365 

366 if isinstance(within, (int, float)): 

367 data["within"] = within = timedelta(seconds=within) 

368 

369 if posture == Posture.Proactive: 

370 if not within or within == timedelta(0): 

371 data["within"] = timedelta(seconds=10.0) 

372 elif within < timedelta(seconds=10.0): 

373 raise ValueError( 

374 "`within` for Proactive triggers must be greater than or equal to " 

375 "10 seconds" 

376 ) 

377 

378 return data 

379 

380 def covers(self, event: ReceivedEvent) -> bool: 

381 if not self.event_pattern.match(event.event): 381 ↛ 382line 381 didn't jump to line 382 because the condition on line 381 was never true

382 return False 

383 

384 if not self.covers_resources(event.resource, event.related): 

385 return False 

386 

387 return True 

388 

389 @property 

390 def immediate(self) -> bool: 

391 """Does this reactive trigger fire immediately for all events?""" 

392 return self.posture == Posture.Reactive and self.within == timedelta(0) 

393 

394 _event_pattern: Optional[re.Pattern[str]] = PrivateAttr(None) 

395 

396 @property 

397 def event_pattern(self) -> re.Pattern[str]: 

398 """A regular expression which may be evaluated against any event string to 

399 determine if this trigger would be interested in the event""" 

400 if self._event_pattern: 

401 return self._event_pattern 

402 

403 if not self.expect: 403 ↛ 408line 403 didn't jump to line 408 because the condition on line 403 was always true

404 # This preserves the trivial match for `expect`, and matches the behavior 

405 # of expects() below 

406 self._event_pattern = re.compile(".+") 

407 else: 

408 patterns = [ 

409 # escape each pattern, then translate wildcards ('*' -> r'.+') 

410 re.escape(e).replace("\\*", ".+") 

411 for e in self.expect | self.after 

412 ] 

413 self._event_pattern = re.compile("|".join(patterns)) 

414 

415 return self._event_pattern 

416 

417 def starts_after(self, event: str) -> bool: 

418 # Warning: Previously we returned 'True' if there was trivial 'after' criteria. 

419 # Although this is not wrong, it led to automations processing more events 

420 # than they should have. 

421 if not self.after: 421 ↛ 422line 421 didn't jump to line 422 because the condition on line 421 was never true

422 return False 

423 

424 for candidate in self.after: 

425 if matches(candidate, event): 425 ↛ 426line 425 didn't jump to line 426 because the condition on line 425 was never true

426 return True 

427 return False 

428 

429 def expects(self, event: str) -> bool: 

430 if not self.expect: 430 ↛ 433line 430 didn't jump to line 433 because the condition on line 430 was always true

431 return True 

432 

433 for candidate in self.expect: 

434 if matches(candidate, event): 

435 return True 

436 return False 

437 

438 def bucketing_key(self, event: ReceivedEvent) -> Tuple[str, ...]: 

439 return tuple( 

440 event.find_resource_label(label) or "" for label in sorted(self.for_each) 

441 ) 

442 

443 def meets_threshold(self, event_count: int) -> bool: 

444 if self.posture == Posture.Reactive and event_count >= self.threshold: 

445 return True 

446 

447 if self.posture == Posture.Proactive and event_count < self.threshold: 447 ↛ 450line 447 didn't jump to line 450 because the condition on line 447 was always true

448 return True 

449 

450 return False 

451 

452 def create_automation_state_change_event( 

453 self, firing: Firing, trigger_state: TriggerState 

454 ) -> ReceivedEvent: 

455 """Returns a ReceivedEvent for an automation state change 

456 into a triggered or resolved state.""" 

457 automation = firing.trigger.automation 

458 triggering_event = firing.triggering_event 

459 

460 resource_data = Resource( 

461 { 

462 "prefect.resource.id": f"prefect.automation.{automation.id}", 

463 "prefect.resource.name": automation.name, 

464 } 

465 ) 

466 

467 if self.posture.value: 467 ↛ 470line 467 didn't jump to line 470 because the condition on line 467 was always true

468 resource_data["prefect.posture"] = self.posture.value 

469 

470 return ReceivedEvent( 

471 occurred=firing.triggered, 

472 event=f"prefect.automation.{trigger_state.value.lower()}", 

473 resource=resource_data, 

474 related=( 

475 [ 

476 RelatedResource( 

477 { 

478 "prefect.resource.id": f"prefect.event.{triggering_event.id}", 

479 "prefect.resource.role": "triggering-event", 

480 } 

481 ) 

482 ] 

483 if triggering_event 

484 else [] 

485 ), 

486 payload={ 

487 "triggering_labels": firing.triggering_labels, 

488 "triggering_event": ( 

489 triggering_event.model_dump(mode="json") 

490 if triggering_event 

491 else None 

492 ), 

493 }, 

494 id=uuid7(), 

495 ) 

496 

497 

498ServerTriggerTypes: TypeAlias = Union[EventTrigger, CompoundTrigger, SequenceTrigger] 

499"""The union of all concrete trigger types that a user may actually create""" 

500 

501T = TypeVar("T", bound=Trigger) 

502 

503 

504class AutomationCore(PrefectBaseModel, extra="ignore"): 

505 """Defines an action a user wants to take when a certain number of events 

506 do or don't happen to the matching resources""" 

507 

508 name: str = Field(default=..., description="The name of this automation") 

509 description: str = Field( 

510 default="", description="A longer description of this automation" 

511 ) 

512 

513 enabled: bool = Field( 

514 default=True, description="Whether this automation will be evaluated" 

515 ) 

516 tags: list[str] = Field( 

517 default_factory=list, 

518 description="A list of tags associated with this automation", 

519 ) 

520 

521 trigger: ServerTriggerTypes = Field( 

522 default=..., 

523 description=( 

524 "The criteria for which events this Automation covers and how it will " 

525 "respond to the presence or absence of those events" 

526 ), 

527 ) 

528 

529 actions: list[ServerActionTypes] = Field( 

530 default=..., 

531 description="The actions to perform when this Automation triggers", 

532 ) 

533 

534 actions_on_trigger: list[ServerActionTypes] = Field( 

535 default_factory=list, 

536 description="The actions to perform when an Automation goes into a triggered state", 

537 ) 

538 

539 actions_on_resolve: list[ServerActionTypes] = Field( 

540 default_factory=list, 

541 description="The actions to perform when an Automation goes into a resolving state", 

542 ) 

543 

544 def triggers(self) -> Sequence[Trigger]: 

545 """Returns all triggers within this automation""" 

546 return self.trigger.all_triggers() 

547 

548 def triggers_of_type(self, trigger_type: Type[T]) -> Sequence[T]: 

549 """Returns all triggers of the specified type within this automation""" 

550 return [t for t in self.triggers() if isinstance(t, trigger_type)] 

551 

552 def trigger_by_id(self, trigger_id: UUID) -> Optional[Trigger]: 

553 """Returns the trigger with the given ID, or None if no such trigger exists""" 

554 for trigger in self.triggers(): 

555 if trigger.id == trigger_id: 

556 return trigger 

557 return None 

558 

559 @model_validator(mode="after") 

560 def prevent_run_deployment_loops(self) -> Self: 

561 """Detects potential infinite loops in automations with RunDeployment actions""" 

562 from prefect.server.events.actions import RunDeployment 

563 

564 if not self.enabled: 

565 # Disabled automations can't cause problems 

566 return self 

567 

568 if ( 

569 not self.trigger 

570 or not isinstance(self.trigger, EventTrigger) 

571 or self.trigger.posture != Posture.Reactive 

572 ): 

573 # Only reactive automations can cause infinite amplification 

574 return self 

575 

576 if not any(e.startswith("prefect.flow-run.") for e in self.trigger.expect): 576 ↛ 581line 576 didn't jump to line 581 because the condition on line 576 was always true

577 # Only flow run events can cause infinite amplification 

578 return self 

579 

580 # Every flow run created by a Deployment goes through these states 

581 problematic_events = { 

582 "prefect.flow-run.Scheduled", 

583 "prefect.flow-run.Pending", 

584 "prefect.flow-run.Running", 

585 "prefect.flow-run.*", 

586 } 

587 if not problematic_events.intersection(self.trigger.expect): 

588 return self 

589 

590 actions = [a for a in self.actions if isinstance(a, RunDeployment)] 

591 for action in actions: 

592 if action.source == "inferred": 

593 # Inferred deployments for flow run state change events will always 

594 # cause infinite loops, because no matter what filters we place on the 

595 # flow run, we're inferring the deployment from it, so we'll always 

596 # produce a new flow run that matches those filters. 

597 raise ValueError( 

598 "Running an inferred deployment from a flow run state change event " 

599 "will lead to an infinite loop of flow runs. Please choose a " 

600 "specific deployment and add additional filtering labels to the " 

601 "match or match_related for this automation's trigger." 

602 ) 

603 

604 if action.source == "selected": 

605 # Selected deployments for flow run state changes can cause infinite 

606 # loops if there aren't enough filtering labels on the trigger's match 

607 # or match_related. While it's still possible to have infinite loops 

608 # with additional filters, it's less likely. 

609 if self.trigger.match.matches_every_resource_of_kind( 

610 "prefect.flow-run" 

611 ): 

612 relateds = ( 

613 self.trigger.match_related 

614 if isinstance(self.trigger.match_related, list) 

615 else [self.trigger.match_related] 

616 ) 

617 if any( 

618 related.matches_every_resource_of_kind("prefect.flow-run") 

619 for related in relateds 

620 ): 

621 raise ValueError( 

622 "Running a selected deployment from a flow run state " 

623 "change event may lead to an infinite loop of flow runs. " 

624 "Please include additional filtering labels on either " 

625 "match or match_related to narrow down which flow runs " 

626 "will trigger this automation to exclude flow runs from " 

627 "the deployment you've selected." 

628 ) 

629 

630 return self 

631 

632 

633class Automation(ORMBaseModel, AutomationCore, extra="ignore"): 

634 def __init__(self, *args: Any, **kwargs: Any): 

635 super().__init__(*args, **kwargs) 

636 self.trigger._set_parent(self) 

637 

638 @classmethod 

639 def model_validate( 

640 cls: type[Self], 

641 obj: Any, 

642 *, 

643 strict: bool | None = None, 

644 from_attributes: bool | None = None, 

645 context: dict[str, Any] | None = None, 

646 ) -> Self: 

647 automation = super().model_validate( 

648 obj, strict=strict, from_attributes=from_attributes, context=context 

649 ) 

650 automation.trigger._set_parent(automation) 

651 return automation 

652 

653 

654class AutomationCreate(AutomationCore, ActionBaseModel, extra="forbid"): 

655 owner_resource: Optional[str] = Field( 

656 default=None, description="The resource to which this automation belongs" 

657 ) 

658 

659 

660class AutomationUpdate(AutomationCore, ActionBaseModel, extra="forbid"): 

661 pass 

662 

663 

664class AutomationPartialUpdate(ActionBaseModel, extra="forbid"): 

665 enabled: bool = Field(True, description="Whether this automation will be evaluated") 

666 

667 

668class AutomationSort(AutoEnum): 

669 """Defines automations sorting options.""" 

670 

671 CREATED_DESC = "CREATED_DESC" 

672 UPDATED_DESC = "UPDATED_DESC" 

673 NAME_ASC = "NAME_ASC" 

674 NAME_DESC = "NAME_DESC" 

675 

676 

677class Firing(PrefectBaseModel): 

678 """Represents one instance of a trigger firing""" 

679 

680 id: UUID = Field(default_factory=uuid7) 

681 

682 trigger: ServerTriggerTypes = Field( 

683 default=..., description="The trigger that is firing" 

684 ) 

685 trigger_states: Set[TriggerState] = Field( 

686 default=..., 

687 description="The state changes represented by this Firing", 

688 ) 

689 triggered: DateTime = Field( 

690 default=..., 

691 description=( 

692 "The time at which this trigger fired, which may differ from the " 

693 "occurred time of the associated event (as events processing may always " 

694 "be slightly delayed)." 

695 ), 

696 ) 

697 triggering_labels: Dict[str, str] = Field( 

698 default_factory=dict, 

699 description=( 

700 "The labels associated with this Firing, derived from the underlying " 

701 "for_each values of the trigger. Only used in the context " 

702 "of EventTriggers." 

703 ), 

704 ) 

705 triggering_firings: List[Firing] = Field( 

706 default_factory=list, 

707 description=( 

708 "The firings of the triggers that caused this trigger to fire. Only used " 

709 "in the context of CompoundTriggers." 

710 ), 

711 ) 

712 triggering_event: Optional[ReceivedEvent] = Field( 

713 default=None, 

714 description=( 

715 "The most recent event associated with this Firing. This may be the " 

716 "event that caused the trigger to fire (for Reactive triggers), or the " 

717 "last event to match the trigger (for Proactive triggers), or the state " 

718 "change event (for a Metric trigger)." 

719 ), 

720 ) 

721 triggering_value: Optional[Any] = Field( 

722 default=None, 

723 description=( 

724 "A value associated with this firing of a trigger. Maybe used to " 

725 "convey additional information at the point of firing, like the value of " 

726 "the last query for a MetricTrigger" 

727 ), 

728 ) 

729 

730 @field_validator("trigger_states") 

731 @classmethod 

732 def validate_trigger_states(cls, value: set[TriggerState]) -> set[TriggerState]: 

733 if not value: 733 ↛ 734line 733 didn't jump to line 734 because the condition on line 733 was never true

734 raise ValueError("At least one trigger state must be provided") 

735 return value 

736 

737 def all_firings(self) -> Sequence[Firing]: 

738 return [self] + [ 

739 f for child in self.triggering_firings for f in child.all_firings() 

740 ] 

741 

742 def all_events(self) -> Sequence[ReceivedEvent]: 

743 events = [self.triggering_event] if self.triggering_event else [] 

744 return events + [ 

745 e for child in self.triggering_firings for e in child.all_events() 

746 ] 

747 

748 

749class TriggeredAction(PrefectBaseModel): 

750 """An action caused as the result of an automation""" 

751 

752 automation: Automation = Field( 

753 ..., description="The Automation that caused this action" 

754 ) 

755 

756 id: UUID = Field( 

757 default_factory=uuid7, 

758 description="A unique key representing a single triggering of an action", 

759 ) 

760 

761 firing: Optional[Firing] = Field( 

762 default=None, description="The Firing that prompted this action" 

763 ) 

764 

765 triggered: DateTime = Field(..., description="When this action was triggered") 

766 triggering_labels: Dict[str, str] = Field( 

767 ..., 

768 description=( 

769 "The subset of labels of the Event that triggered this action, " 

770 "corresponding to the Automation's for_each. If no for_each is specified, " 

771 "this will be an empty set of labels" 

772 ), 

773 ) 

774 triggering_event: Optional[ReceivedEvent] = Field( 

775 ..., 

776 description=( 

777 "The last Event to trigger this automation, if applicable. For reactive " 

778 "triggers, this will be the event that caused the trigger to fire. For " 

779 "proactive triggers, this will be the last event to match the automation, " 

780 "if there was one." 

781 ), 

782 ) 

783 action: ServerActionTypes = Field( 

784 ..., 

785 description="The action to perform", 

786 ) 

787 action_index: int = Field( 

788 default=0, 

789 description="The index of the action within the automation", 

790 ) 

791 automation_triggered_event_id: UUID | None = Field( 

792 default=None, 

793 description=( 

794 "The ID of the automation.triggered or automation.resolved event that " 

795 "prompted this action, used to link automation.action.* events back to " 

796 "the state change event" 

797 ), 

798 ) 

799 

800 def idempotency_key(self) -> str: 

801 """Produce a human-friendly idempotency key for this action""" 

802 return ", ".join( 

803 [ 

804 f"automation {self.automation.id}", 

805 f"action {self.action_index}", 

806 f"invocation {self.id}", 

807 ] 

808 ) 

809 

810 def all_firings(self) -> Sequence[Firing]: 

811 return self.firing.all_firings() if self.firing else [] 

812 

813 def all_events(self) -> Sequence[ReceivedEvent]: 

814 return self.firing.all_events() if self.firing else [] 

815 

816 

817CompoundTrigger.model_rebuild() 

818SequenceTrigger.model_rebuild()