queue.py 16 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609
  1. # -*- coding: utf-8 -*-
  2. """
  3. Asynchrone Event-Queue.
  4. Queue für Hintergrund-Verarbeitung von Events.
  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. import uuid
  13. if TYPE_CHECKING:
  14. from trixy_core.events.event_data.base import EventData
  15. class QueuePriority(Enum):
  16. """
  17. Prioritäten für Queue-Einträge.
  18. """
  19. LOW = 30
  20. NORMAL = 20
  21. HIGH = 10
  22. CRITICAL = 0
  23. @dataclass
  24. class QueuedEvent:
  25. """
  26. Ein Event in der Queue.
  27. """
  28. id: str
  29. """Eindeutige ID."""
  30. event_name: str
  31. """Name des Events."""
  32. event_data: "EventData"
  33. """Die Event-Daten."""
  34. priority: QueuePriority = QueuePriority.NORMAL
  35. """Priorität."""
  36. created_at: datetime = field(default_factory=datetime.now)
  37. """Zeitpunkt der Erstellung."""
  38. retry_count: int = 0
  39. """Anzahl der Wiederholungen."""
  40. metadata: dict[str, Any] = field(default_factory=dict)
  41. """Zusätzliche Metadaten."""
  42. @classmethod
  43. def create(
  44. cls,
  45. event_name: str,
  46. event_data: "EventData",
  47. priority: QueuePriority = QueuePriority.NORMAL,
  48. ) -> "QueuedEvent":
  49. """
  50. Erstellt einen neuen Queue-Eintrag.
  51. Args:
  52. event_name: Name des Events.
  53. event_data: Die Event-Daten.
  54. priority: Priorität.
  55. Returns:
  56. QueuedEvent.
  57. """
  58. return cls(
  59. id=str(uuid.uuid4()),
  60. event_name=event_name,
  61. event_data=event_data,
  62. priority=priority,
  63. )
  64. def __lt__(self, other: "QueuedEvent") -> bool:
  65. """Vergleich für Priority-Queue."""
  66. if self.priority.value != other.priority.value:
  67. return self.priority.value < other.priority.value
  68. return self.created_at < other.created_at
  69. @dataclass
  70. class QueueConfig:
  71. """
  72. Konfiguration für die Event-Queue.
  73. """
  74. max_size: int = 10000
  75. """Maximale Queue-Größe."""
  76. default_priority: QueuePriority = QueuePriority.NORMAL
  77. """Standard-Priorität."""
  78. enable_priority: bool = True
  79. """Prioritäts-Queue aktivieren."""
  80. overflow_policy: str = "drop_oldest"
  81. """Overflow-Verhalten: drop_oldest, drop_newest, block."""
  82. max_retries: int = 3
  83. """Maximale Wiederholungen."""
  84. @dataclass
  85. class QueueStats:
  86. """
  87. Statistiken der Event-Queue.
  88. """
  89. total_enqueued: int = 0
  90. """Gesamtzahl eingestellter Events."""
  91. total_processed: int = 0
  92. """Gesamtzahl verarbeiteter Events."""
  93. total_failed: int = 0
  94. """Gesamtzahl fehlgeschlagener Events."""
  95. total_dropped: int = 0
  96. """Gesamtzahl verworfener Events."""
  97. current_size: int = 0
  98. """Aktuelle Queue-Größe."""
  99. by_priority: dict[str, int] = field(default_factory=dict)
  100. """Events nach Priorität."""
  101. avg_wait_time_ms: float = 0.0
  102. """Durchschnittliche Wartezeit."""
  103. def to_dict(self) -> dict[str, Any]:
  104. """Konvertiert in ein Dictionary."""
  105. return {
  106. "total_enqueued": self.total_enqueued,
  107. "total_processed": self.total_processed,
  108. "total_failed": self.total_failed,
  109. "total_dropped": self.total_dropped,
  110. "current_size": self.current_size,
  111. "by_priority": self.by_priority,
  112. "avg_wait_time_ms": self.avg_wait_time_ms,
  113. }
  114. class AsyncEventQueue:
  115. """
  116. Asynchrone Event-Queue.
  117. Ermöglicht die Hintergrund-Verarbeitung von Events
  118. mit Prioritäten und Wiederholungen.
  119. Example:
  120. queue = AsyncEventQueue(QueueConfig(max_size=1000))
  121. # Handler registrieren
  122. queue.set_handler(process_event)
  123. # Events einstellen
  124. await queue.enqueue("user_login", event_data)
  125. await queue.enqueue("critical_alert", data, QueuePriority.CRITICAL)
  126. # Verarbeitung starten
  127. await queue.start()
  128. # Später stoppen
  129. await queue.stop()
  130. """
  131. def __init__(
  132. self,
  133. config: QueueConfig | None = None,
  134. ) -> None:
  135. """
  136. Initialisiert die Queue.
  137. Args:
  138. config: Konfiguration.
  139. """
  140. self._config = config or QueueConfig()
  141. if self._config.enable_priority:
  142. self._queue: asyncio.PriorityQueue[QueuedEvent] = asyncio.PriorityQueue(
  143. maxsize=self._config.max_size
  144. )
  145. else:
  146. self._queue = asyncio.Queue(maxsize=self._config.max_size)
  147. self._handler: Callable[[str, "EventData"], Coroutine[Any, Any, None]] | None = None
  148. self._running = False
  149. self._workers: list[asyncio.Task] = []
  150. self._stats = QueueStats()
  151. self._wait_times: list[float] = []
  152. @property
  153. def size(self) -> int:
  154. """Aktuelle Queue-Größe."""
  155. return self._queue.qsize()
  156. @property
  157. def is_full(self) -> bool:
  158. """Ob Queue voll ist."""
  159. return self._queue.full()
  160. @property
  161. def is_empty(self) -> bool:
  162. """Ob Queue leer ist."""
  163. return self._queue.empty()
  164. @property
  165. def is_running(self) -> bool:
  166. """Ob Queue läuft."""
  167. return self._running
  168. def set_handler(
  169. self,
  170. handler: Callable[[str, "EventData"], Coroutine[Any, Any, None]],
  171. ) -> None:
  172. """
  173. Setzt den Event-Handler.
  174. Args:
  175. handler: Handler-Funktion.
  176. """
  177. self._handler = handler
  178. async def enqueue(
  179. self,
  180. event_name: str,
  181. event_data: "EventData",
  182. priority: QueuePriority | None = None,
  183. ) -> bool:
  184. """
  185. Fügt ein Event zur Queue hinzu.
  186. Args:
  187. event_name: Name des Events.
  188. event_data: Die Event-Daten.
  189. priority: Optionale Priorität.
  190. Returns:
  191. True wenn erfolgreich.
  192. """
  193. priority = priority or self._config.default_priority
  194. queued = QueuedEvent.create(event_name, event_data, priority)
  195. try:
  196. if self._queue.full():
  197. if self._config.overflow_policy == "drop_oldest":
  198. try:
  199. self._queue.get_nowait()
  200. self._stats.total_dropped += 1
  201. except asyncio.QueueEmpty:
  202. pass
  203. elif self._config.overflow_policy == "drop_newest":
  204. self._stats.total_dropped += 1
  205. return False
  206. elif self._config.overflow_policy == "block":
  207. await self._queue.put(queued)
  208. self._stats.total_enqueued += 1
  209. return True
  210. self._queue.put_nowait(queued)
  211. self._stats.total_enqueued += 1
  212. # Statistiken nach Priorität
  213. prio_name = priority.name
  214. self._stats.by_priority[prio_name] = self._stats.by_priority.get(prio_name, 0) + 1
  215. return True
  216. except asyncio.QueueFull:
  217. self._stats.total_dropped += 1
  218. return False
  219. async def dequeue(self, timeout: float | None = None) -> QueuedEvent | None:
  220. """
  221. Holt das nächste Event aus der Queue.
  222. Args:
  223. timeout: Optionales Timeout in Sekunden.
  224. Returns:
  225. QueuedEvent oder None bei Timeout.
  226. """
  227. try:
  228. if timeout is not None:
  229. queued = await asyncio.wait_for(self._queue.get(), timeout=timeout)
  230. else:
  231. queued = await self._queue.get()
  232. # Wartezeit berechnen
  233. wait_time_ms = (datetime.now() - queued.created_at).total_seconds() * 1000
  234. self._wait_times.append(wait_time_ms)
  235. if len(self._wait_times) > 100:
  236. self._wait_times = self._wait_times[-50:]
  237. self._stats.avg_wait_time_ms = sum(self._wait_times) / len(self._wait_times)
  238. return queued
  239. except asyncio.TimeoutError:
  240. return None
  241. async def peek(self) -> QueuedEvent | None:
  242. """
  243. Gibt das nächste Event zurück ohne zu entfernen.
  244. Returns:
  245. QueuedEvent oder None.
  246. """
  247. # Für Priority-Queue nicht direkt möglich
  248. # Workaround: Get und sofort wieder Put
  249. try:
  250. queued = self._queue.get_nowait()
  251. await self._queue.put(queued)
  252. return queued
  253. except asyncio.QueueEmpty:
  254. return None
  255. async def process_one(self) -> bool:
  256. """
  257. Verarbeitet ein Event.
  258. Returns:
  259. True wenn erfolgreich.
  260. """
  261. if not self._handler:
  262. raise RuntimeError("Kein Handler registriert")
  263. queued = await self.dequeue(timeout=1.0)
  264. if not queued:
  265. return False
  266. try:
  267. await self._handler(queued.event_name, queued.event_data)
  268. self._stats.total_processed += 1
  269. return True
  270. except Exception as e:
  271. queued.retry_count += 1
  272. if queued.retry_count < self._config.max_retries:
  273. # Zurück in die Queue
  274. await self._queue.put(queued)
  275. else:
  276. self._stats.total_failed += 1
  277. return False
  278. async def start(self, num_workers: int = 1) -> None:
  279. """
  280. Startet die Queue-Verarbeitung.
  281. Args:
  282. num_workers: Anzahl der Worker.
  283. """
  284. if self._running:
  285. return
  286. if not self._handler:
  287. raise RuntimeError("Kein Handler registriert")
  288. self._running = True
  289. async def worker():
  290. while self._running:
  291. try:
  292. await self.process_one()
  293. except Exception:
  294. await asyncio.sleep(0.1)
  295. for _ in range(num_workers):
  296. task = asyncio.create_task(worker())
  297. self._workers.append(task)
  298. async def stop(self, wait: bool = True, timeout: float = 10.0) -> None:
  299. """
  300. Stoppt die Queue-Verarbeitung.
  301. Args:
  302. wait: Ob auf Abschluss gewartet wird.
  303. timeout: Timeout in Sekunden.
  304. """
  305. self._running = False
  306. if wait and self._workers:
  307. try:
  308. await asyncio.wait_for(
  309. asyncio.gather(*self._workers, return_exceptions=True),
  310. timeout=timeout,
  311. )
  312. except asyncio.TimeoutError:
  313. for worker in self._workers:
  314. worker.cancel()
  315. self._workers.clear()
  316. async def drain(self) -> int:
  317. """
  318. Verarbeitet alle wartenden Events.
  319. Returns:
  320. Anzahl verarbeiteter Events.
  321. """
  322. if not self._handler:
  323. raise RuntimeError("Kein Handler registriert")
  324. count = 0
  325. while not self._queue.empty():
  326. if await self.process_one():
  327. count += 1
  328. return count
  329. def clear(self) -> int:
  330. """
  331. Löscht alle wartenden Events.
  332. Returns:
  333. Anzahl gelöschter Events.
  334. """
  335. count = 0
  336. while not self._queue.empty():
  337. try:
  338. self._queue.get_nowait()
  339. count += 1
  340. except asyncio.QueueEmpty:
  341. break
  342. return count
  343. def get_stats(self) -> QueueStats:
  344. """
  345. Gibt Statistiken zurück.
  346. Returns:
  347. QueueStats.
  348. """
  349. self._stats.current_size = self.size
  350. return self._stats
  351. class MultiQueue:
  352. """
  353. Multi-Queue mit separaten Queues pro Event-Typ.
  354. Ermöglicht unterschiedliche Konfigurationen
  355. für verschiedene Event-Typen.
  356. Example:
  357. multi = MultiQueue()
  358. multi.create_queue("high_priority", QueueConfig(max_size=100))
  359. multi.create_queue("bulk", QueueConfig(max_size=10000))
  360. multi.route("critical_*", "high_priority")
  361. multi.route("*", "bulk")
  362. await multi.enqueue("critical_alert", data) # -> high_priority
  363. await multi.enqueue("user_click", data) # -> bulk
  364. """
  365. def __init__(self) -> None:
  366. """Initialisiert die Multi-Queue."""
  367. self._queues: dict[str, AsyncEventQueue] = {}
  368. self._routes: list[tuple[str, str]] = [] # (pattern, queue_name)
  369. self._default_queue: str | None = None
  370. def create_queue(
  371. self,
  372. name: str,
  373. config: QueueConfig | None = None,
  374. ) -> AsyncEventQueue:
  375. """
  376. Erstellt eine Queue.
  377. Args:
  378. name: Name der Queue.
  379. config: Konfiguration.
  380. Returns:
  381. Die erstellte Queue.
  382. """
  383. queue = AsyncEventQueue(config)
  384. self._queues[name] = queue
  385. return queue
  386. def get_queue(self, name: str) -> AsyncEventQueue | None:
  387. """
  388. Gibt eine Queue zurück.
  389. Args:
  390. name: Name der Queue.
  391. Returns:
  392. AsyncEventQueue oder None.
  393. """
  394. return self._queues.get(name)
  395. def set_default_queue(self, name: str) -> None:
  396. """
  397. Setzt die Standard-Queue.
  398. Args:
  399. name: Name der Queue.
  400. """
  401. if name not in self._queues:
  402. raise ValueError(f"Queue '{name}' existiert nicht")
  403. self._default_queue = name
  404. def route(self, pattern: str, queue_name: str) -> None:
  405. """
  406. Fügt eine Routing-Regel hinzu.
  407. Args:
  408. pattern: Event-Name-Pattern (mit * als Wildcard).
  409. queue_name: Ziel-Queue.
  410. """
  411. if queue_name not in self._queues:
  412. raise ValueError(f"Queue '{queue_name}' existiert nicht")
  413. self._routes.append((pattern, queue_name))
  414. def _match_pattern(self, event_name: str, pattern: str) -> bool:
  415. """Prüft ob Event-Name zum Pattern passt."""
  416. import fnmatch
  417. return fnmatch.fnmatch(event_name, pattern)
  418. def _get_target_queue(self, event_name: str) -> str | None:
  419. """Ermittelt die Ziel-Queue für ein Event."""
  420. for pattern, queue_name in self._routes:
  421. if self._match_pattern(event_name, pattern):
  422. return queue_name
  423. return self._default_queue
  424. async def enqueue(
  425. self,
  426. event_name: str,
  427. event_data: "EventData",
  428. priority: QueuePriority | None = None,
  429. queue_name: str | None = None,
  430. ) -> bool:
  431. """
  432. Fügt ein Event zur passenden Queue hinzu.
  433. Args:
  434. event_name: Name des Events.
  435. event_data: Die Event-Daten.
  436. priority: Optionale Priorität.
  437. queue_name: Optionaler Queue-Name (überschreibt Routing).
  438. Returns:
  439. True wenn erfolgreich.
  440. """
  441. target = queue_name or self._get_target_queue(event_name)
  442. if not target or target not in self._queues:
  443. return False
  444. return await self._queues[target].enqueue(event_name, event_data, priority)
  445. def set_handler(
  446. self,
  447. queue_name: str,
  448. handler: Callable[[str, "EventData"], Coroutine[Any, Any, None]],
  449. ) -> None:
  450. """
  451. Setzt den Handler für eine Queue.
  452. Args:
  453. queue_name: Name der Queue.
  454. handler: Handler-Funktion.
  455. """
  456. if queue_name not in self._queues:
  457. raise ValueError(f"Queue '{queue_name}' existiert nicht")
  458. self._queues[queue_name].set_handler(handler)
  459. async def start_all(self, num_workers_per_queue: int = 1) -> None:
  460. """
  461. Startet alle Queues.
  462. Args:
  463. num_workers_per_queue: Worker pro Queue.
  464. """
  465. for queue in self._queues.values():
  466. if queue._handler:
  467. await queue.start(num_workers_per_queue)
  468. async def stop_all(self) -> None:
  469. """Stoppt alle Queues."""
  470. for queue in self._queues.values():
  471. await queue.stop()
  472. def get_all_stats(self) -> dict[str, QueueStats]:
  473. """
  474. Gibt Statistiken aller Queues zurück.
  475. Returns:
  476. Dictionary von Queue-Name zu Stats.
  477. """
  478. return {
  479. name: queue.get_stats()
  480. for name, queue in self._queues.items()
  481. }