replay.py 12 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434
  1. # -*- coding: utf-8 -*-
  2. """
  3. Event Replayer für Wiedergabe gespeicherter Events.
  4. Ermöglicht die Wiedergabe von Events zu Test- und Debug-Zwecken.
  5. """
  6. from __future__ import annotations
  7. import asyncio
  8. from dataclasses import dataclass, field
  9. from datetime import datetime
  10. from enum import Enum, auto
  11. from typing import TYPE_CHECKING, Any, Callable, Coroutine
  12. from trixy_core.events.persistence.store import (
  13. EventStore,
  14. StoredEvent,
  15. EventQuery,
  16. )
  17. if TYPE_CHECKING:
  18. from trixy_core.events.event_data.base import EventData
  19. class ReplayMode(Enum):
  20. """
  21. Wiedergabe-Modi.
  22. """
  23. REALTIME = auto()
  24. """Echtzeit-Wiedergabe mit originalen Zeitabständen."""
  25. FAST = auto()
  26. """So schnell wie möglich."""
  27. SCALED = auto()
  28. """Skalierte Geschwindigkeit."""
  29. STEPPED = auto()
  30. """Schrittweise (wartet auf Signal)."""
  31. @dataclass
  32. class ReplayConfig:
  33. """
  34. Konfiguration für Event-Wiedergabe.
  35. """
  36. mode: ReplayMode = ReplayMode.FAST
  37. """Wiedergabe-Modus."""
  38. speed_factor: float = 1.0
  39. """Geschwindigkeitsfaktor für SCALED-Modus."""
  40. batch_size: int = 100
  41. """Batch-Größe für Verarbeitung."""
  42. pause_between_batches_ms: float = 0.0
  43. """Pause zwischen Batches."""
  44. filter_events: list[str] | None = None
  45. """Events die wiedergegeben werden."""
  46. exclude_events: list[str] | None = None
  47. """Events die ausgeschlossen werden."""
  48. start_time: datetime | None = None
  49. """Start-Zeitpunkt für Wiedergabe."""
  50. end_time: datetime | None = None
  51. """End-Zeitpunkt für Wiedergabe."""
  52. @dataclass
  53. class ReplayResult:
  54. """
  55. Ergebnis einer Wiedergabe.
  56. """
  57. total_events: int = 0
  58. """Gesamtzahl der Events."""
  59. replayed_events: int = 0
  60. """Anzahl wiedergegebener Events."""
  61. skipped_events: int = 0
  62. """Anzahl übersprungener Events."""
  63. failed_events: int = 0
  64. """Anzahl fehlgeschlagener Events."""
  65. duration_ms: float = 0.0
  66. """Dauer in Millisekunden."""
  67. started_at: datetime | None = None
  68. """Startzeitpunkt."""
  69. completed_at: datetime | None = None
  70. """Endzeitpunkt."""
  71. errors: list[str] = field(default_factory=list)
  72. """Aufgetretene Fehler."""
  73. @property
  74. def events_per_second(self) -> float:
  75. """Wiedergegebene Events pro Sekunde."""
  76. if self.duration_ms <= 0:
  77. return 0.0
  78. return self.replayed_events / (self.duration_ms / 1000)
  79. def to_dict(self) -> dict[str, Any]:
  80. """Konvertiert in ein Dictionary."""
  81. return {
  82. "total_events": self.total_events,
  83. "replayed_events": self.replayed_events,
  84. "skipped_events": self.skipped_events,
  85. "failed_events": self.failed_events,
  86. "duration_ms": self.duration_ms,
  87. "events_per_second": self.events_per_second,
  88. "started_at": self.started_at.isoformat() if self.started_at else None,
  89. "completed_at": self.completed_at.isoformat() if self.completed_at else None,
  90. "errors": self.errors,
  91. }
  92. class EventReplayer:
  93. """
  94. Replayer für gespeicherte Events.
  95. Wiederholt Events aus einem Store für Tests,
  96. Debugging oder Analyse.
  97. Example:
  98. store = InMemoryEventStore()
  99. # ... Events speichern ...
  100. replayer = EventReplayer(store)
  101. # Handler registrieren
  102. replayer.on_event("user_login", handle_login)
  103. replayer.on_event("*", log_all_events)
  104. # Wiedergabe starten
  105. result = await replayer.replay(ReplayConfig(
  106. mode=ReplayMode.FAST,
  107. filter_events=["user_login", "user_logout"],
  108. ))
  109. print(f"Replayed {result.replayed_events} events")
  110. """
  111. def __init__(
  112. self,
  113. store: EventStore,
  114. ) -> None:
  115. """
  116. Initialisiert den Replayer.
  117. Args:
  118. store: Der Event Store.
  119. """
  120. self._store = store
  121. self._handlers: dict[str, list[Callable[[StoredEvent], Coroutine[Any, Any, None]]]] = {}
  122. self._global_handlers: list[Callable[[StoredEvent], Coroutine[Any, Any, None]]] = []
  123. self._running = False
  124. self._paused = False
  125. self._step_event = asyncio.Event()
  126. self._cancel_event = asyncio.Event()
  127. @property
  128. def is_running(self) -> bool:
  129. """Ob Wiedergabe läuft."""
  130. return self._running
  131. @property
  132. def is_paused(self) -> bool:
  133. """Ob Wiedergabe pausiert."""
  134. return self._paused
  135. def on_event(
  136. self,
  137. event_name: str,
  138. handler: Callable[[StoredEvent], Coroutine[Any, Any, None]],
  139. ) -> None:
  140. """
  141. Registriert einen Handler für ein Event.
  142. Args:
  143. event_name: Event-Name oder "*" für alle.
  144. handler: Handler-Funktion.
  145. """
  146. if event_name == "*":
  147. self._global_handlers.append(handler)
  148. else:
  149. if event_name not in self._handlers:
  150. self._handlers[event_name] = []
  151. self._handlers[event_name].append(handler)
  152. def remove_handler(
  153. self,
  154. event_name: str,
  155. handler: Callable[[StoredEvent], Coroutine[Any, Any, None]],
  156. ) -> bool:
  157. """
  158. Entfernt einen Handler.
  159. Args:
  160. event_name: Event-Name.
  161. handler: Handler-Funktion.
  162. Returns:
  163. True wenn entfernt.
  164. """
  165. if event_name == "*":
  166. if handler in self._global_handlers:
  167. self._global_handlers.remove(handler)
  168. return True
  169. else:
  170. handlers = self._handlers.get(event_name, [])
  171. if handler in handlers:
  172. handlers.remove(handler)
  173. return True
  174. return False
  175. def clear_handlers(self) -> None:
  176. """Entfernt alle Handler."""
  177. self._handlers.clear()
  178. self._global_handlers.clear()
  179. async def replay(
  180. self,
  181. config: ReplayConfig | None = None,
  182. ) -> ReplayResult:
  183. """
  184. Startet die Wiedergabe.
  185. Args:
  186. config: Wiedergabe-Konfiguration.
  187. Returns:
  188. ReplayResult.
  189. """
  190. if self._running:
  191. raise RuntimeError("Wiedergabe läuft bereits")
  192. config = config or ReplayConfig()
  193. result = ReplayResult()
  194. result.started_at = datetime.now()
  195. self._running = True
  196. self._paused = False
  197. self._cancel_event.clear()
  198. try:
  199. # Query erstellen
  200. query = EventQuery(
  201. event_names=config.filter_events,
  202. start_time=config.start_time,
  203. end_time=config.end_time,
  204. order_ascending=True,
  205. )
  206. # Events laden
  207. query_result = await self._store.query(query)
  208. events = query_result.events
  209. result.total_events = len(events)
  210. last_timestamp: datetime | None = None
  211. for i, event in enumerate(events):
  212. # Abbruch prüfen
  213. if self._cancel_event.is_set():
  214. break
  215. # Pause prüfen
  216. while self._paused and not self._cancel_event.is_set():
  217. if config.mode == ReplayMode.STEPPED:
  218. self._step_event.clear()
  219. await self._step_event.wait()
  220. break
  221. await asyncio.sleep(0.1)
  222. # Exclude-Filter
  223. if config.exclude_events and event.event_name in config.exclude_events:
  224. result.skipped_events += 1
  225. continue
  226. # Timing berechnen
  227. if config.mode == ReplayMode.REALTIME and last_timestamp:
  228. delay = (event.timestamp - last_timestamp).total_seconds()
  229. if delay > 0:
  230. await asyncio.sleep(delay)
  231. elif config.mode == ReplayMode.SCALED and last_timestamp:
  232. delay = (event.timestamp - last_timestamp).total_seconds()
  233. if delay > 0:
  234. await asyncio.sleep(delay / config.speed_factor)
  235. last_timestamp = event.timestamp
  236. # Handler aufrufen
  237. try:
  238. await self._invoke_handlers(event)
  239. result.replayed_events += 1
  240. except Exception as e:
  241. result.failed_events += 1
  242. result.errors.append(f"{event.event_name}: {e}")
  243. # Batch-Pause
  244. if config.pause_between_batches_ms > 0:
  245. if (i + 1) % config.batch_size == 0:
  246. await asyncio.sleep(config.pause_between_batches_ms / 1000)
  247. finally:
  248. self._running = False
  249. result.completed_at = datetime.now()
  250. result.duration_ms = (
  251. result.completed_at - result.started_at
  252. ).total_seconds() * 1000
  253. return result
  254. async def _invoke_handlers(self, event: StoredEvent) -> None:
  255. """
  256. Ruft Handler für ein Event auf.
  257. Args:
  258. event: Das Event.
  259. """
  260. # Spezifische Handler
  261. handlers = self._handlers.get(event.event_name, [])
  262. for handler in handlers:
  263. await handler(event)
  264. # Globale Handler
  265. for handler in self._global_handlers:
  266. await handler(event)
  267. def pause(self) -> None:
  268. """Pausiert die Wiedergabe."""
  269. self._paused = True
  270. def resume(self) -> None:
  271. """Setzt die Wiedergabe fort."""
  272. self._paused = False
  273. self._step_event.set()
  274. def step(self) -> None:
  275. """Führt einen Schritt aus (für STEPPED-Modus)."""
  276. self._step_event.set()
  277. def cancel(self) -> None:
  278. """Bricht die Wiedergabe ab."""
  279. self._cancel_event.set()
  280. self._step_event.set() # Falls wartend
  281. async def replay_single(
  282. self,
  283. event_id: str,
  284. ) -> bool:
  285. """
  286. Wiederholt ein einzelnes Event.
  287. Args:
  288. event_id: Die Event-ID.
  289. Returns:
  290. True wenn erfolgreich.
  291. """
  292. event = await self._store.get(event_id)
  293. if not event:
  294. return False
  295. try:
  296. await self._invoke_handlers(event)
  297. return True
  298. except Exception:
  299. return False
  300. async def replay_range(
  301. self,
  302. start_sequence: int,
  303. end_sequence: int,
  304. config: ReplayConfig | None = None,
  305. ) -> ReplayResult:
  306. """
  307. Wiederholt Events in einem Sequenzbereich.
  308. Args:
  309. start_sequence: Start-Sequenz.
  310. end_sequence: End-Sequenz.
  311. config: Wiedergabe-Konfiguration.
  312. Returns:
  313. ReplayResult.
  314. """
  315. config = config or ReplayConfig()
  316. # Query mit Sequenz-Filter
  317. query = EventQuery(
  318. event_names=config.filter_events,
  319. after_sequence=start_sequence - 1,
  320. before_sequence=end_sequence + 1,
  321. order_ascending=True,
  322. )
  323. query_result = await self._store.query(query)
  324. # Temporär Events setzen und replay aufrufen
  325. result = ReplayResult()
  326. result.started_at = datetime.now()
  327. result.total_events = len(query_result.events)
  328. for event in query_result.events:
  329. if config.exclude_events and event.event_name in config.exclude_events:
  330. result.skipped_events += 1
  331. continue
  332. try:
  333. await self._invoke_handlers(event)
  334. result.replayed_events += 1
  335. except Exception as e:
  336. result.failed_events += 1
  337. result.errors.append(f"{event.event_name}: {e}")
  338. result.completed_at = datetime.now()
  339. result.duration_ms = (
  340. result.completed_at - result.started_at
  341. ).total_seconds() * 1000
  342. return result