event.py 9.5 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289
  1. # -*- coding: utf-8 -*-
  2. """
  3. Event Action - Emittiert Events über EventManager.
  4. """
  5. from typing import Any, Callable
  6. from trixy_core.scheduler.action.base import Action
  7. from trixy_core.scheduler.form_fields import FieldType, FormField, FormValidationError
  8. class EmitEventAction(Action):
  9. """
  10. Action die ein Event über den EventManager emittiert.
  11. Beispiel:
  12. - EmitEventAction("reminder_triggered", {"message": "Zeit für eine Pause!"})
  13. - EmitEventAction("daily_report", data_fn=lambda ctx: {"date": ctx["date"]})
  14. """
  15. def __init__(
  16. self,
  17. event_name: str,
  18. event_data: dict | None = None,
  19. data_fn: Callable[[dict], dict] | None = None,
  20. name: str | None = None,
  21. timeout_seconds: float = 10.0,
  22. ):
  23. """
  24. Initialisiert die Action.
  25. Args:
  26. event_name: Name des Events
  27. event_data: Statische Event-Daten
  28. data_fn: Funktion für dynamische Event-Daten (context → data)
  29. name: Action-Name
  30. timeout_seconds: Timeout
  31. """
  32. super().__init__(name or f"emit_{event_name}", timeout_seconds)
  33. self._event_name = event_name
  34. self._event_data = event_data or {}
  35. self._data_fn = data_fn
  36. @property
  37. def event_name(self) -> str:
  38. """Event-Name."""
  39. return self._event_name
  40. @property
  41. def event_data(self) -> dict:
  42. """Statische Event-Daten."""
  43. return dict(self._event_data)
  44. async def _execute(self, context: dict) -> dict:
  45. """Emittiert das Event."""
  46. # Event-Manager aus Kontext holen
  47. event_manager = context.get("event_manager")
  48. # Daten zusammenstellen
  49. data = dict(self._event_data)
  50. # Dynamische Daten
  51. if self._data_fn:
  52. dynamic_data = self._data_fn(context)
  53. data.update(dynamic_data)
  54. # Job-Info hinzufügen
  55. if "job_id" in context:
  56. data["_source_job"] = context["job_id"]
  57. # Event emittieren
  58. if event_manager:
  59. await event_manager.emit(self._event_name, data)
  60. else:
  61. # Fallback: Event als Ergebnis zurückgeben
  62. pass
  63. return {
  64. "event_name": self._event_name,
  65. "event_data": data,
  66. }
  67. def to_dict(self) -> dict:
  68. """Konvertiert zu Dictionary."""
  69. return {
  70. **self._base_dict(),
  71. "event_name": self._event_name,
  72. "event_data": self._event_data,
  73. # data_fn kann nicht serialisiert werden
  74. }
  75. @classmethod
  76. def from_dict(cls, data: dict) -> "EmitEventAction":
  77. """Erstellt Action aus Dictionary."""
  78. return cls(
  79. event_name=data["event_name"],
  80. event_data=data.get("event_data", {}),
  81. name=data.get("name"),
  82. timeout_seconds=data.get("timeout_seconds", 10.0),
  83. )
  84. @classmethod
  85. def display_name(cls) -> str:
  86. """Anzeigename fuer die TUI."""
  87. return "Event senden"
  88. @classmethod
  89. def form_fields(cls) -> list[FormField]:
  90. """Gibt die Formular-Felder fuer die TUI zurueck."""
  91. return [
  92. FormField(
  93. name="event_name", label="Event-Name", required=True,
  94. placeholder="reminder_triggered",
  95. ),
  96. FormField(
  97. name="event_data", label="Event-Daten (JSON)",
  98. field_type=FieldType.TEXTAREA, lines=4,
  99. placeholder='{"key": "value"}',
  100. help_text="JSON-Objekt mit Event-Daten",
  101. ),
  102. FormField(
  103. name="timeout_seconds", label="Timeout (Sek.)",
  104. field_type=FieldType.NUMBER,
  105. default=10.0, min_value=1.0, is_float=True,
  106. ),
  107. ]
  108. @classmethod
  109. def validate_form(cls, data: dict) -> list[FormValidationError]:
  110. """Validiert Formular-Daten. Gibt Liste von Fehlern zurueck."""
  111. errors: list[FormValidationError] = []
  112. if not data.get("event_name", "").strip():
  113. errors.append(FormValidationError("event_name", "Event-Name ist erforderlich"))
  114. event_data = data.get("event_data", "")
  115. if event_data and isinstance(event_data, str):
  116. try:
  117. import json
  118. parsed = json.loads(event_data)
  119. if not isinstance(parsed, dict):
  120. errors.append(FormValidationError("event_data", "Muss ein JSON-Objekt sein"))
  121. except json.JSONDecodeError as e:
  122. errors.append(FormValidationError("event_data", f"Ungueltiges JSON: {e}"))
  123. return errors
  124. class MultiEventAction(Action):
  125. """
  126. Action die mehrere Events emittiert.
  127. Beispiel:
  128. - MultiEventAction([
  129. ("event_a", {"type": "a"}),
  130. ("event_b", {"type": "b"}),
  131. ])
  132. """
  133. def __init__(
  134. self,
  135. events: list[tuple[str, dict]],
  136. sequential: bool = True,
  137. name: str | None = None,
  138. timeout_seconds: float = 30.0,
  139. ):
  140. """
  141. Initialisiert die Action.
  142. Args:
  143. events: Liste von (event_name, event_data) Tupeln
  144. sequential: Events nacheinander (True) oder parallel (False)
  145. name: Action-Name
  146. timeout_seconds: Timeout
  147. """
  148. super().__init__(name or "multi_event", timeout_seconds)
  149. self._events = list(events)
  150. self._sequential = sequential
  151. @property
  152. def events(self) -> list[tuple[str, dict]]:
  153. """Event-Liste."""
  154. return list(self._events)
  155. def add_event(self, event_name: str, event_data: dict | None = None) -> None:
  156. """Fügt Event hinzu."""
  157. self._events.append((event_name, event_data or {}))
  158. async def _execute(self, context: dict) -> dict:
  159. """Emittiert alle Events."""
  160. event_manager = context.get("event_manager")
  161. results = []
  162. if self._sequential:
  163. # Nacheinander
  164. for event_name, event_data in self._events:
  165. data = dict(event_data)
  166. if "job_id" in context:
  167. data["_source_job"] = context["job_id"]
  168. if event_manager:
  169. await event_manager.emit(event_name, data)
  170. results.append({
  171. "event_name": event_name,
  172. "event_data": data,
  173. })
  174. else:
  175. # Parallel
  176. import asyncio
  177. async def emit_event(name: str, data: dict) -> dict:
  178. d = dict(data)
  179. if "job_id" in context:
  180. d["_source_job"] = context["job_id"]
  181. if event_manager:
  182. await event_manager.emit(name, d)
  183. return {"event_name": name, "event_data": d}
  184. tasks = [
  185. emit_event(name, data)
  186. for name, data in self._events
  187. ]
  188. results = await asyncio.gather(*tasks)
  189. return {
  190. "events_emitted": len(results),
  191. "results": list(results),
  192. }
  193. def to_dict(self) -> dict:
  194. """Konvertiert zu Dictionary."""
  195. return {
  196. **self._base_dict(),
  197. "events": self._events,
  198. "sequential": self._sequential,
  199. }
  200. @classmethod
  201. def from_dict(cls, data: dict) -> "MultiEventAction":
  202. """Erstellt Action aus Dictionary."""
  203. return cls(
  204. events=[(e[0], e[1]) for e in data["events"]],
  205. sequential=data.get("sequential", True),
  206. name=data.get("name"),
  207. timeout_seconds=data.get("timeout_seconds", 30.0),
  208. )
  209. @classmethod
  210. def display_name(cls) -> str:
  211. """Anzeigename fuer die TUI."""
  212. return "Mehrere Events senden"
  213. @classmethod
  214. def form_fields(cls) -> list[FormField]:
  215. """Gibt die Formular-Felder fuer die TUI zurueck."""
  216. return [
  217. FormField(
  218. name="events", label="Events (JSON)",
  219. field_type=FieldType.TEXTAREA, required=True, lines=6,
  220. placeholder='[["event_a", {"key": "val"}], ["event_b", {}]]',
  221. help_text="JSON-Array von [event_name, event_data] Paaren",
  222. ),
  223. FormField(
  224. name="sequential", label="Sequentiell ausfuehren",
  225. field_type=FieldType.CHECKBOX, default=True,
  226. help_text="Events nacheinander (True) oder parallel (False) senden",
  227. ),
  228. ]
  229. @classmethod
  230. def validate_form(cls, data: dict) -> list[FormValidationError]:
  231. """Validiert Formular-Daten. Gibt Liste von Fehlern zurueck."""
  232. errors: list[FormValidationError] = []
  233. events_raw = data.get("events", "")
  234. if isinstance(events_raw, str):
  235. if not events_raw.strip():
  236. errors.append(FormValidationError("events", "Events-Liste ist erforderlich"))
  237. else:
  238. try:
  239. import json
  240. parsed = json.loads(events_raw)
  241. if not isinstance(parsed, list):
  242. errors.append(FormValidationError("events", "Muss ein JSON-Array sein"))
  243. elif not parsed:
  244. errors.append(FormValidationError("events", "Mindestens ein Event erforderlich"))
  245. except json.JSONDecodeError as e:
  246. errors.append(FormValidationError("events", f"Ungueltiges JSON: {e}"))
  247. elif isinstance(events_raw, list):
  248. if not events_raw:
  249. errors.append(FormValidationError("events", "Mindestens ein Event erforderlich"))
  250. return errors