| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609 |
- # -*- coding: utf-8 -*-
- """
- Asynchrone Event-Queue.
- Queue für Hintergrund-Verarbeitung von Events.
- """
- 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
- import uuid
- if TYPE_CHECKING:
- from trixy_core.events.event_data.base import EventData
- class QueuePriority(Enum):
- """
- Prioritäten für Queue-Einträge.
- """
- LOW = 30
- NORMAL = 20
- HIGH = 10
- CRITICAL = 0
- @dataclass
- class QueuedEvent:
- """
- Ein Event in der Queue.
- """
- id: str
- """Eindeutige ID."""
- event_name: str
- """Name des Events."""
- event_data: "EventData"
- """Die Event-Daten."""
- priority: QueuePriority = QueuePriority.NORMAL
- """Priorität."""
- created_at: datetime = field(default_factory=datetime.now)
- """Zeitpunkt der Erstellung."""
- retry_count: int = 0
- """Anzahl der Wiederholungen."""
- metadata: dict[str, Any] = field(default_factory=dict)
- """Zusätzliche Metadaten."""
- @classmethod
- def create(
- cls,
- event_name: str,
- event_data: "EventData",
- priority: QueuePriority = QueuePriority.NORMAL,
- ) -> "QueuedEvent":
- """
- Erstellt einen neuen Queue-Eintrag.
- Args:
- event_name: Name des Events.
- event_data: Die Event-Daten.
- priority: Priorität.
- Returns:
- QueuedEvent.
- """
- return cls(
- id=str(uuid.uuid4()),
- event_name=event_name,
- event_data=event_data,
- priority=priority,
- )
- def __lt__(self, other: "QueuedEvent") -> bool:
- """Vergleich für Priority-Queue."""
- if self.priority.value != other.priority.value:
- return self.priority.value < other.priority.value
- return self.created_at < other.created_at
- @dataclass
- class QueueConfig:
- """
- Konfiguration für die Event-Queue.
- """
- max_size: int = 10000
- """Maximale Queue-Größe."""
- default_priority: QueuePriority = QueuePriority.NORMAL
- """Standard-Priorität."""
- enable_priority: bool = True
- """Prioritäts-Queue aktivieren."""
- overflow_policy: str = "drop_oldest"
- """Overflow-Verhalten: drop_oldest, drop_newest, block."""
- max_retries: int = 3
- """Maximale Wiederholungen."""
- @dataclass
- class QueueStats:
- """
- Statistiken der Event-Queue.
- """
- total_enqueued: int = 0
- """Gesamtzahl eingestellter Events."""
- total_processed: int = 0
- """Gesamtzahl verarbeiteter Events."""
- total_failed: int = 0
- """Gesamtzahl fehlgeschlagener Events."""
- total_dropped: int = 0
- """Gesamtzahl verworfener Events."""
- current_size: int = 0
- """Aktuelle Queue-Größe."""
- by_priority: dict[str, int] = field(default_factory=dict)
- """Events nach Priorität."""
- avg_wait_time_ms: float = 0.0
- """Durchschnittliche Wartezeit."""
- def to_dict(self) -> dict[str, Any]:
- """Konvertiert in ein Dictionary."""
- return {
- "total_enqueued": self.total_enqueued,
- "total_processed": self.total_processed,
- "total_failed": self.total_failed,
- "total_dropped": self.total_dropped,
- "current_size": self.current_size,
- "by_priority": self.by_priority,
- "avg_wait_time_ms": self.avg_wait_time_ms,
- }
- class AsyncEventQueue:
- """
- Asynchrone Event-Queue.
- Ermöglicht die Hintergrund-Verarbeitung von Events
- mit Prioritäten und Wiederholungen.
- Example:
- queue = AsyncEventQueue(QueueConfig(max_size=1000))
- # Handler registrieren
- queue.set_handler(process_event)
- # Events einstellen
- await queue.enqueue("user_login", event_data)
- await queue.enqueue("critical_alert", data, QueuePriority.CRITICAL)
- # Verarbeitung starten
- await queue.start()
- # Später stoppen
- await queue.stop()
- """
- def __init__(
- self,
- config: QueueConfig | None = None,
- ) -> None:
- """
- Initialisiert die Queue.
- Args:
- config: Konfiguration.
- """
- self._config = config or QueueConfig()
- if self._config.enable_priority:
- self._queue: asyncio.PriorityQueue[QueuedEvent] = asyncio.PriorityQueue(
- maxsize=self._config.max_size
- )
- else:
- self._queue = asyncio.Queue(maxsize=self._config.max_size)
- self._handler: Callable[[str, "EventData"], Coroutine[Any, Any, None]] | None = None
- self._running = False
- self._workers: list[asyncio.Task] = []
- self._stats = QueueStats()
- self._wait_times: list[float] = []
- @property
- def size(self) -> int:
- """Aktuelle Queue-Größe."""
- return self._queue.qsize()
- @property
- def is_full(self) -> bool:
- """Ob Queue voll ist."""
- return self._queue.full()
- @property
- def is_empty(self) -> bool:
- """Ob Queue leer ist."""
- return self._queue.empty()
- @property
- def is_running(self) -> bool:
- """Ob Queue läuft."""
- return self._running
- def set_handler(
- self,
- handler: Callable[[str, "EventData"], Coroutine[Any, Any, None]],
- ) -> None:
- """
- Setzt den Event-Handler.
- Args:
- handler: Handler-Funktion.
- """
- self._handler = handler
- async def enqueue(
- self,
- event_name: str,
- event_data: "EventData",
- priority: QueuePriority | None = None,
- ) -> bool:
- """
- Fügt ein Event zur Queue hinzu.
- Args:
- event_name: Name des Events.
- event_data: Die Event-Daten.
- priority: Optionale Priorität.
- Returns:
- True wenn erfolgreich.
- """
- priority = priority or self._config.default_priority
- queued = QueuedEvent.create(event_name, event_data, priority)
- try:
- if self._queue.full():
- if self._config.overflow_policy == "drop_oldest":
- try:
- self._queue.get_nowait()
- self._stats.total_dropped += 1
- except asyncio.QueueEmpty:
- pass
- elif self._config.overflow_policy == "drop_newest":
- self._stats.total_dropped += 1
- return False
- elif self._config.overflow_policy == "block":
- await self._queue.put(queued)
- self._stats.total_enqueued += 1
- return True
- self._queue.put_nowait(queued)
- self._stats.total_enqueued += 1
- # Statistiken nach Priorität
- prio_name = priority.name
- self._stats.by_priority[prio_name] = self._stats.by_priority.get(prio_name, 0) + 1
- return True
- except asyncio.QueueFull:
- self._stats.total_dropped += 1
- return False
- async def dequeue(self, timeout: float | None = None) -> QueuedEvent | None:
- """
- Holt das nächste Event aus der Queue.
- Args:
- timeout: Optionales Timeout in Sekunden.
- Returns:
- QueuedEvent oder None bei Timeout.
- """
- try:
- if timeout is not None:
- queued = await asyncio.wait_for(self._queue.get(), timeout=timeout)
- else:
- queued = await self._queue.get()
- # Wartezeit berechnen
- wait_time_ms = (datetime.now() - queued.created_at).total_seconds() * 1000
- self._wait_times.append(wait_time_ms)
- if len(self._wait_times) > 100:
- self._wait_times = self._wait_times[-50:]
- self._stats.avg_wait_time_ms = sum(self._wait_times) / len(self._wait_times)
- return queued
- except asyncio.TimeoutError:
- return None
- async def peek(self) -> QueuedEvent | None:
- """
- Gibt das nächste Event zurück ohne zu entfernen.
- Returns:
- QueuedEvent oder None.
- """
- # Für Priority-Queue nicht direkt möglich
- # Workaround: Get und sofort wieder Put
- try:
- queued = self._queue.get_nowait()
- await self._queue.put(queued)
- return queued
- except asyncio.QueueEmpty:
- return None
- async def process_one(self) -> bool:
- """
- Verarbeitet ein Event.
- Returns:
- True wenn erfolgreich.
- """
- if not self._handler:
- raise RuntimeError("Kein Handler registriert")
- queued = await self.dequeue(timeout=1.0)
- if not queued:
- return False
- try:
- await self._handler(queued.event_name, queued.event_data)
- self._stats.total_processed += 1
- return True
- except Exception as e:
- queued.retry_count += 1
- if queued.retry_count < self._config.max_retries:
- # Zurück in die Queue
- await self._queue.put(queued)
- else:
- self._stats.total_failed += 1
- return False
- async def start(self, num_workers: int = 1) -> None:
- """
- Startet die Queue-Verarbeitung.
- Args:
- num_workers: Anzahl der Worker.
- """
- if self._running:
- return
- if not self._handler:
- raise RuntimeError("Kein Handler registriert")
- self._running = True
- async def worker():
- while self._running:
- try:
- await self.process_one()
- except Exception:
- await asyncio.sleep(0.1)
- for _ in range(num_workers):
- task = asyncio.create_task(worker())
- self._workers.append(task)
- async def stop(self, wait: bool = True, timeout: float = 10.0) -> None:
- """
- Stoppt die Queue-Verarbeitung.
- Args:
- wait: Ob auf Abschluss gewartet wird.
- timeout: Timeout in Sekunden.
- """
- self._running = False
- if wait and self._workers:
- try:
- await asyncio.wait_for(
- asyncio.gather(*self._workers, return_exceptions=True),
- timeout=timeout,
- )
- except asyncio.TimeoutError:
- for worker in self._workers:
- worker.cancel()
- self._workers.clear()
- async def drain(self) -> int:
- """
- Verarbeitet alle wartenden Events.
- Returns:
- Anzahl verarbeiteter Events.
- """
- if not self._handler:
- raise RuntimeError("Kein Handler registriert")
- count = 0
- while not self._queue.empty():
- if await self.process_one():
- count += 1
- return count
- def clear(self) -> int:
- """
- Löscht alle wartenden Events.
- Returns:
- Anzahl gelöschter Events.
- """
- count = 0
- while not self._queue.empty():
- try:
- self._queue.get_nowait()
- count += 1
- except asyncio.QueueEmpty:
- break
- return count
- def get_stats(self) -> QueueStats:
- """
- Gibt Statistiken zurück.
- Returns:
- QueueStats.
- """
- self._stats.current_size = self.size
- return self._stats
- class MultiQueue:
- """
- Multi-Queue mit separaten Queues pro Event-Typ.
- Ermöglicht unterschiedliche Konfigurationen
- für verschiedene Event-Typen.
- Example:
- multi = MultiQueue()
- multi.create_queue("high_priority", QueueConfig(max_size=100))
- multi.create_queue("bulk", QueueConfig(max_size=10000))
- multi.route("critical_*", "high_priority")
- multi.route("*", "bulk")
- await multi.enqueue("critical_alert", data) # -> high_priority
- await multi.enqueue("user_click", data) # -> bulk
- """
- def __init__(self) -> None:
- """Initialisiert die Multi-Queue."""
- self._queues: dict[str, AsyncEventQueue] = {}
- self._routes: list[tuple[str, str]] = [] # (pattern, queue_name)
- self._default_queue: str | None = None
- def create_queue(
- self,
- name: str,
- config: QueueConfig | None = None,
- ) -> AsyncEventQueue:
- """
- Erstellt eine Queue.
- Args:
- name: Name der Queue.
- config: Konfiguration.
- Returns:
- Die erstellte Queue.
- """
- queue = AsyncEventQueue(config)
- self._queues[name] = queue
- return queue
- def get_queue(self, name: str) -> AsyncEventQueue | None:
- """
- Gibt eine Queue zurück.
- Args:
- name: Name der Queue.
- Returns:
- AsyncEventQueue oder None.
- """
- return self._queues.get(name)
- def set_default_queue(self, name: str) -> None:
- """
- Setzt die Standard-Queue.
- Args:
- name: Name der Queue.
- """
- if name not in self._queues:
- raise ValueError(f"Queue '{name}' existiert nicht")
- self._default_queue = name
- def route(self, pattern: str, queue_name: str) -> None:
- """
- Fügt eine Routing-Regel hinzu.
- Args:
- pattern: Event-Name-Pattern (mit * als Wildcard).
- queue_name: Ziel-Queue.
- """
- if queue_name not in self._queues:
- raise ValueError(f"Queue '{queue_name}' existiert nicht")
- self._routes.append((pattern, queue_name))
- def _match_pattern(self, event_name: str, pattern: str) -> bool:
- """Prüft ob Event-Name zum Pattern passt."""
- import fnmatch
- return fnmatch.fnmatch(event_name, pattern)
- def _get_target_queue(self, event_name: str) -> str | None:
- """Ermittelt die Ziel-Queue für ein Event."""
- for pattern, queue_name in self._routes:
- if self._match_pattern(event_name, pattern):
- return queue_name
- return self._default_queue
- async def enqueue(
- self,
- event_name: str,
- event_data: "EventData",
- priority: QueuePriority | None = None,
- queue_name: str | None = None,
- ) -> bool:
- """
- Fügt ein Event zur passenden Queue hinzu.
- Args:
- event_name: Name des Events.
- event_data: Die Event-Daten.
- priority: Optionale Priorität.
- queue_name: Optionaler Queue-Name (überschreibt Routing).
- Returns:
- True wenn erfolgreich.
- """
- target = queue_name or self._get_target_queue(event_name)
- if not target or target not in self._queues:
- return False
- return await self._queues[target].enqueue(event_name, event_data, priority)
- def set_handler(
- self,
- queue_name: str,
- handler: Callable[[str, "EventData"], Coroutine[Any, Any, None]],
- ) -> None:
- """
- Setzt den Handler für eine Queue.
- Args:
- queue_name: Name der Queue.
- handler: Handler-Funktion.
- """
- if queue_name not in self._queues:
- raise ValueError(f"Queue '{queue_name}' existiert nicht")
- self._queues[queue_name].set_handler(handler)
- async def start_all(self, num_workers_per_queue: int = 1) -> None:
- """
- Startet alle Queues.
- Args:
- num_workers_per_queue: Worker pro Queue.
- """
- for queue in self._queues.values():
- if queue._handler:
- await queue.start(num_workers_per_queue)
- async def stop_all(self) -> None:
- """Stoppt alle Queues."""
- for queue in self._queues.values():
- await queue.stop()
- def get_all_stats(self) -> dict[str, QueueStats]:
- """
- Gibt Statistiken aller Queues zurück.
- Returns:
- Dictionary von Queue-Name zu Stats.
- """
- return {
- name: queue.get_stats()
- for name, queue in self._queues.items()
- }
|