| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434 |
- # -*- coding: utf-8 -*-
- """
- Event Replayer für Wiedergabe gespeicherter Events.
- Ermöglicht die Wiedergabe von Events zu Test- und Debug-Zwecken.
- """
- from __future__ import annotations
- import asyncio
- from dataclasses import dataclass, field
- from datetime import datetime
- from enum import Enum, auto
- from typing import TYPE_CHECKING, Any, Callable, Coroutine
- from trixy_core.events.persistence.store import (
- EventStore,
- StoredEvent,
- EventQuery,
- )
- if TYPE_CHECKING:
- from trixy_core.events.event_data.base import EventData
- class ReplayMode(Enum):
- """
- Wiedergabe-Modi.
- """
- REALTIME = auto()
- """Echtzeit-Wiedergabe mit originalen Zeitabständen."""
- FAST = auto()
- """So schnell wie möglich."""
- SCALED = auto()
- """Skalierte Geschwindigkeit."""
- STEPPED = auto()
- """Schrittweise (wartet auf Signal)."""
- @dataclass
- class ReplayConfig:
- """
- Konfiguration für Event-Wiedergabe.
- """
- mode: ReplayMode = ReplayMode.FAST
- """Wiedergabe-Modus."""
- speed_factor: float = 1.0
- """Geschwindigkeitsfaktor für SCALED-Modus."""
- batch_size: int = 100
- """Batch-Größe für Verarbeitung."""
- pause_between_batches_ms: float = 0.0
- """Pause zwischen Batches."""
- filter_events: list[str] | None = None
- """Events die wiedergegeben werden."""
- exclude_events: list[str] | None = None
- """Events die ausgeschlossen werden."""
- start_time: datetime | None = None
- """Start-Zeitpunkt für Wiedergabe."""
- end_time: datetime | None = None
- """End-Zeitpunkt für Wiedergabe."""
- @dataclass
- class ReplayResult:
- """
- Ergebnis einer Wiedergabe.
- """
- total_events: int = 0
- """Gesamtzahl der Events."""
- replayed_events: int = 0
- """Anzahl wiedergegebener Events."""
- skipped_events: int = 0
- """Anzahl übersprungener Events."""
- failed_events: int = 0
- """Anzahl fehlgeschlagener Events."""
- duration_ms: float = 0.0
- """Dauer in Millisekunden."""
- started_at: datetime | None = None
- """Startzeitpunkt."""
- completed_at: datetime | None = None
- """Endzeitpunkt."""
- errors: list[str] = field(default_factory=list)
- """Aufgetretene Fehler."""
- @property
- def events_per_second(self) -> float:
- """Wiedergegebene Events pro Sekunde."""
- if self.duration_ms <= 0:
- return 0.0
- return self.replayed_events / (self.duration_ms / 1000)
- def to_dict(self) -> dict[str, Any]:
- """Konvertiert in ein Dictionary."""
- return {
- "total_events": self.total_events,
- "replayed_events": self.replayed_events,
- "skipped_events": self.skipped_events,
- "failed_events": self.failed_events,
- "duration_ms": self.duration_ms,
- "events_per_second": self.events_per_second,
- "started_at": self.started_at.isoformat() if self.started_at else None,
- "completed_at": self.completed_at.isoformat() if self.completed_at else None,
- "errors": self.errors,
- }
- class EventReplayer:
- """
- Replayer für gespeicherte Events.
- Wiederholt Events aus einem Store für Tests,
- Debugging oder Analyse.
- Example:
- store = InMemoryEventStore()
- # ... Events speichern ...
- replayer = EventReplayer(store)
- # Handler registrieren
- replayer.on_event("user_login", handle_login)
- replayer.on_event("*", log_all_events)
- # Wiedergabe starten
- result = await replayer.replay(ReplayConfig(
- mode=ReplayMode.FAST,
- filter_events=["user_login", "user_logout"],
- ))
- print(f"Replayed {result.replayed_events} events")
- """
- def __init__(
- self,
- store: EventStore,
- ) -> None:
- """
- Initialisiert den Replayer.
- Args:
- store: Der Event Store.
- """
- self._store = store
- self._handlers: dict[str, list[Callable[[StoredEvent], Coroutine[Any, Any, None]]]] = {}
- self._global_handlers: list[Callable[[StoredEvent], Coroutine[Any, Any, None]]] = []
- self._running = False
- self._paused = False
- self._step_event = asyncio.Event()
- self._cancel_event = asyncio.Event()
- @property
- def is_running(self) -> bool:
- """Ob Wiedergabe läuft."""
- return self._running
- @property
- def is_paused(self) -> bool:
- """Ob Wiedergabe pausiert."""
- return self._paused
- def on_event(
- self,
- event_name: str,
- handler: Callable[[StoredEvent], Coroutine[Any, Any, None]],
- ) -> None:
- """
- Registriert einen Handler für ein Event.
- Args:
- event_name: Event-Name oder "*" für alle.
- handler: Handler-Funktion.
- """
- if event_name == "*":
- self._global_handlers.append(handler)
- else:
- if event_name not in self._handlers:
- self._handlers[event_name] = []
- self._handlers[event_name].append(handler)
- def remove_handler(
- self,
- event_name: str,
- handler: Callable[[StoredEvent], Coroutine[Any, Any, None]],
- ) -> bool:
- """
- Entfernt einen Handler.
- Args:
- event_name: Event-Name.
- handler: Handler-Funktion.
- Returns:
- True wenn entfernt.
- """
- if event_name == "*":
- if handler in self._global_handlers:
- self._global_handlers.remove(handler)
- return True
- else:
- handlers = self._handlers.get(event_name, [])
- if handler in handlers:
- handlers.remove(handler)
- return True
- return False
- def clear_handlers(self) -> None:
- """Entfernt alle Handler."""
- self._handlers.clear()
- self._global_handlers.clear()
- async def replay(
- self,
- config: ReplayConfig | None = None,
- ) -> ReplayResult:
- """
- Startet die Wiedergabe.
- Args:
- config: Wiedergabe-Konfiguration.
- Returns:
- ReplayResult.
- """
- if self._running:
- raise RuntimeError("Wiedergabe läuft bereits")
- config = config or ReplayConfig()
- result = ReplayResult()
- result.started_at = datetime.now()
- self._running = True
- self._paused = False
- self._cancel_event.clear()
- try:
- # Query erstellen
- query = EventQuery(
- event_names=config.filter_events,
- start_time=config.start_time,
- end_time=config.end_time,
- order_ascending=True,
- )
- # Events laden
- query_result = await self._store.query(query)
- events = query_result.events
- result.total_events = len(events)
- last_timestamp: datetime | None = None
- for i, event in enumerate(events):
- # Abbruch prüfen
- if self._cancel_event.is_set():
- break
- # Pause prüfen
- while self._paused and not self._cancel_event.is_set():
- if config.mode == ReplayMode.STEPPED:
- self._step_event.clear()
- await self._step_event.wait()
- break
- await asyncio.sleep(0.1)
- # Exclude-Filter
- if config.exclude_events and event.event_name in config.exclude_events:
- result.skipped_events += 1
- continue
- # Timing berechnen
- if config.mode == ReplayMode.REALTIME and last_timestamp:
- delay = (event.timestamp - last_timestamp).total_seconds()
- if delay > 0:
- await asyncio.sleep(delay)
- elif config.mode == ReplayMode.SCALED and last_timestamp:
- delay = (event.timestamp - last_timestamp).total_seconds()
- if delay > 0:
- await asyncio.sleep(delay / config.speed_factor)
- last_timestamp = event.timestamp
- # Handler aufrufen
- try:
- await self._invoke_handlers(event)
- result.replayed_events += 1
- except Exception as e:
- result.failed_events += 1
- result.errors.append(f"{event.event_name}: {e}")
- # Batch-Pause
- if config.pause_between_batches_ms > 0:
- if (i + 1) % config.batch_size == 0:
- await asyncio.sleep(config.pause_between_batches_ms / 1000)
- finally:
- self._running = False
- result.completed_at = datetime.now()
- result.duration_ms = (
- result.completed_at - result.started_at
- ).total_seconds() * 1000
- return result
- async def _invoke_handlers(self, event: StoredEvent) -> None:
- """
- Ruft Handler für ein Event auf.
- Args:
- event: Das Event.
- """
- # Spezifische Handler
- handlers = self._handlers.get(event.event_name, [])
- for handler in handlers:
- await handler(event)
- # Globale Handler
- for handler in self._global_handlers:
- await handler(event)
- def pause(self) -> None:
- """Pausiert die Wiedergabe."""
- self._paused = True
- def resume(self) -> None:
- """Setzt die Wiedergabe fort."""
- self._paused = False
- self._step_event.set()
- def step(self) -> None:
- """Führt einen Schritt aus (für STEPPED-Modus)."""
- self._step_event.set()
- def cancel(self) -> None:
- """Bricht die Wiedergabe ab."""
- self._cancel_event.set()
- self._step_event.set() # Falls wartend
- async def replay_single(
- self,
- event_id: str,
- ) -> bool:
- """
- Wiederholt ein einzelnes Event.
- Args:
- event_id: Die Event-ID.
- Returns:
- True wenn erfolgreich.
- """
- event = await self._store.get(event_id)
- if not event:
- return False
- try:
- await self._invoke_handlers(event)
- return True
- except Exception:
- return False
- async def replay_range(
- self,
- start_sequence: int,
- end_sequence: int,
- config: ReplayConfig | None = None,
- ) -> ReplayResult:
- """
- Wiederholt Events in einem Sequenzbereich.
- Args:
- start_sequence: Start-Sequenz.
- end_sequence: End-Sequenz.
- config: Wiedergabe-Konfiguration.
- Returns:
- ReplayResult.
- """
- config = config or ReplayConfig()
- # Query mit Sequenz-Filter
- query = EventQuery(
- event_names=config.filter_events,
- after_sequence=start_sequence - 1,
- before_sequence=end_sequence + 1,
- order_ascending=True,
- )
- query_result = await self._store.query(query)
- # Temporär Events setzen und replay aufrufen
- result = ReplayResult()
- result.started_at = datetime.now()
- result.total_events = len(query_result.events)
- for event in query_result.events:
- if config.exclude_events and event.event_name in config.exclude_events:
- result.skipped_events += 1
- continue
- try:
- await self._invoke_handlers(event)
- result.replayed_events += 1
- except Exception as e:
- result.failed_events += 1
- result.errors.append(f"{event.event_name}: {e}")
- result.completed_at = datetime.now()
- result.duration_ms = (
- result.completed_at - result.started_at
- ).total_seconds() * 1000
- return result
|