| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289 |
- # -*- coding: utf-8 -*-
- """
- Event Action - Emittiert Events über EventManager.
- """
- from typing import Any, Callable
- from trixy_core.scheduler.action.base import Action
- from trixy_core.scheduler.form_fields import FieldType, FormField, FormValidationError
- class EmitEventAction(Action):
- """
- Action die ein Event über den EventManager emittiert.
- Beispiel:
- - EmitEventAction("reminder_triggered", {"message": "Zeit für eine Pause!"})
- - EmitEventAction("daily_report", data_fn=lambda ctx: {"date": ctx["date"]})
- """
- def __init__(
- self,
- event_name: str,
- event_data: dict | None = None,
- data_fn: Callable[[dict], dict] | None = None,
- name: str | None = None,
- timeout_seconds: float = 10.0,
- ):
- """
- Initialisiert die Action.
- Args:
- event_name: Name des Events
- event_data: Statische Event-Daten
- data_fn: Funktion für dynamische Event-Daten (context → data)
- name: Action-Name
- timeout_seconds: Timeout
- """
- super().__init__(name or f"emit_{event_name}", timeout_seconds)
- self._event_name = event_name
- self._event_data = event_data or {}
- self._data_fn = data_fn
- @property
- def event_name(self) -> str:
- """Event-Name."""
- return self._event_name
- @property
- def event_data(self) -> dict:
- """Statische Event-Daten."""
- return dict(self._event_data)
- async def _execute(self, context: dict) -> dict:
- """Emittiert das Event."""
- # Event-Manager aus Kontext holen
- event_manager = context.get("event_manager")
- # Daten zusammenstellen
- data = dict(self._event_data)
- # Dynamische Daten
- if self._data_fn:
- dynamic_data = self._data_fn(context)
- data.update(dynamic_data)
- # Job-Info hinzufügen
- if "job_id" in context:
- data["_source_job"] = context["job_id"]
- # Event emittieren
- if event_manager:
- await event_manager.emit(self._event_name, data)
- else:
- # Fallback: Event als Ergebnis zurückgeben
- pass
- return {
- "event_name": self._event_name,
- "event_data": data,
- }
- def to_dict(self) -> dict:
- """Konvertiert zu Dictionary."""
- return {
- **self._base_dict(),
- "event_name": self._event_name,
- "event_data": self._event_data,
- # data_fn kann nicht serialisiert werden
- }
- @classmethod
- def from_dict(cls, data: dict) -> "EmitEventAction":
- """Erstellt Action aus Dictionary."""
- return cls(
- event_name=data["event_name"],
- event_data=data.get("event_data", {}),
- name=data.get("name"),
- timeout_seconds=data.get("timeout_seconds", 10.0),
- )
- @classmethod
- def display_name(cls) -> str:
- """Anzeigename fuer die TUI."""
- return "Event senden"
- @classmethod
- def form_fields(cls) -> list[FormField]:
- """Gibt die Formular-Felder fuer die TUI zurueck."""
- return [
- FormField(
- name="event_name", label="Event-Name", required=True,
- placeholder="reminder_triggered",
- ),
- FormField(
- name="event_data", label="Event-Daten (JSON)",
- field_type=FieldType.TEXTAREA, lines=4,
- placeholder='{"key": "value"}',
- help_text="JSON-Objekt mit Event-Daten",
- ),
- FormField(
- name="timeout_seconds", label="Timeout (Sek.)",
- field_type=FieldType.NUMBER,
- default=10.0, min_value=1.0, is_float=True,
- ),
- ]
- @classmethod
- def validate_form(cls, data: dict) -> list[FormValidationError]:
- """Validiert Formular-Daten. Gibt Liste von Fehlern zurueck."""
- errors: list[FormValidationError] = []
- if not data.get("event_name", "").strip():
- errors.append(FormValidationError("event_name", "Event-Name ist erforderlich"))
- event_data = data.get("event_data", "")
- if event_data and isinstance(event_data, str):
- try:
- import json
- parsed = json.loads(event_data)
- if not isinstance(parsed, dict):
- errors.append(FormValidationError("event_data", "Muss ein JSON-Objekt sein"))
- except json.JSONDecodeError as e:
- errors.append(FormValidationError("event_data", f"Ungueltiges JSON: {e}"))
- return errors
- class MultiEventAction(Action):
- """
- Action die mehrere Events emittiert.
- Beispiel:
- - MultiEventAction([
- ("event_a", {"type": "a"}),
- ("event_b", {"type": "b"}),
- ])
- """
- def __init__(
- self,
- events: list[tuple[str, dict]],
- sequential: bool = True,
- name: str | None = None,
- timeout_seconds: float = 30.0,
- ):
- """
- Initialisiert die Action.
- Args:
- events: Liste von (event_name, event_data) Tupeln
- sequential: Events nacheinander (True) oder parallel (False)
- name: Action-Name
- timeout_seconds: Timeout
- """
- super().__init__(name or "multi_event", timeout_seconds)
- self._events = list(events)
- self._sequential = sequential
- @property
- def events(self) -> list[tuple[str, dict]]:
- """Event-Liste."""
- return list(self._events)
- def add_event(self, event_name: str, event_data: dict | None = None) -> None:
- """Fügt Event hinzu."""
- self._events.append((event_name, event_data or {}))
- async def _execute(self, context: dict) -> dict:
- """Emittiert alle Events."""
- event_manager = context.get("event_manager")
- results = []
- if self._sequential:
- # Nacheinander
- for event_name, event_data in self._events:
- data = dict(event_data)
- if "job_id" in context:
- data["_source_job"] = context["job_id"]
- if event_manager:
- await event_manager.emit(event_name, data)
- results.append({
- "event_name": event_name,
- "event_data": data,
- })
- else:
- # Parallel
- import asyncio
- async def emit_event(name: str, data: dict) -> dict:
- d = dict(data)
- if "job_id" in context:
- d["_source_job"] = context["job_id"]
- if event_manager:
- await event_manager.emit(name, d)
- return {"event_name": name, "event_data": d}
- tasks = [
- emit_event(name, data)
- for name, data in self._events
- ]
- results = await asyncio.gather(*tasks)
- return {
- "events_emitted": len(results),
- "results": list(results),
- }
- def to_dict(self) -> dict:
- """Konvertiert zu Dictionary."""
- return {
- **self._base_dict(),
- "events": self._events,
- "sequential": self._sequential,
- }
- @classmethod
- def from_dict(cls, data: dict) -> "MultiEventAction":
- """Erstellt Action aus Dictionary."""
- return cls(
- events=[(e[0], e[1]) for e in data["events"]],
- sequential=data.get("sequential", True),
- name=data.get("name"),
- timeout_seconds=data.get("timeout_seconds", 30.0),
- )
- @classmethod
- def display_name(cls) -> str:
- """Anzeigename fuer die TUI."""
- return "Mehrere Events senden"
- @classmethod
- def form_fields(cls) -> list[FormField]:
- """Gibt die Formular-Felder fuer die TUI zurueck."""
- return [
- FormField(
- name="events", label="Events (JSON)",
- field_type=FieldType.TEXTAREA, required=True, lines=6,
- placeholder='[["event_a", {"key": "val"}], ["event_b", {}]]',
- help_text="JSON-Array von [event_name, event_data] Paaren",
- ),
- FormField(
- name="sequential", label="Sequentiell ausfuehren",
- field_type=FieldType.CHECKBOX, default=True,
- help_text="Events nacheinander (True) oder parallel (False) senden",
- ),
- ]
- @classmethod
- def validate_form(cls, data: dict) -> list[FormValidationError]:
- """Validiert Formular-Daten. Gibt Liste von Fehlern zurueck."""
- errors: list[FormValidationError] = []
- events_raw = data.get("events", "")
- if isinstance(events_raw, str):
- if not events_raw.strip():
- errors.append(FormValidationError("events", "Events-Liste ist erforderlich"))
- else:
- try:
- import json
- parsed = json.loads(events_raw)
- if not isinstance(parsed, list):
- errors.append(FormValidationError("events", "Muss ein JSON-Array sein"))
- elif not parsed:
- errors.append(FormValidationError("events", "Mindestens ein Event erforderlich"))
- except json.JSONDecodeError as e:
- errors.append(FormValidationError("events", f"Ungueltiges JSON: {e}"))
- elif isinstance(events_raw, list):
- if not events_raw:
- errors.append(FormValidationError("events", "Mindestens ein Event erforderlich"))
- return errors
|