| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430 |
- # -*- coding: utf-8 -*-
- """
- Connection Pool für effiziente Verbindungsverwaltung.
- Bietet Pooling mit Größenlimits, Timeouts und Health-Checks.
- """
- import asyncio
- import ssl
- from collections import deque
- from dataclasses import dataclass
- from typing import Any, Callable, Awaitable
- from trixy_core.network.pool.connection import Connection, ConnectionState
- @dataclass
- class PoolConfig:
- """
- Konfiguration für den Connection Pool.
- Attributes:
- min_size: Minimale Pool-Größe
- max_size: Maximale Pool-Größe
- max_idle_time: Maximale Idle-Zeit in Sekunden
- max_lifetime: Maximale Lebensdauer einer Verbindung
- connect_timeout: Timeout für Verbindungsaufbau
- acquire_timeout: Timeout für Pool-Acquire
- health_check_interval: Intervall für Health-Checks
- validate_on_acquire: Verbindung vor Ausgabe prüfen
- ssl_context: Optionaler SSL-Kontext
- """
- min_size: int = 1
- max_size: int = 10
- max_idle_time: float = 300.0 # 5 Minuten
- max_lifetime: float = 3600.0 # 1 Stunde
- connect_timeout: float = 10.0
- acquire_timeout: float = 30.0
- health_check_interval: float = 30.0
- validate_on_acquire: bool = True
- ssl_context: ssl.SSLContext | None = None
- @dataclass
- class PoolStatistics:
- """Statistiken des Connection Pools."""
- total_connections: int = 0
- idle_connections: int = 0
- in_use_connections: int = 0
- total_acquires: int = 0
- total_releases: int = 0
- total_timeouts: int = 0
- total_errors: int = 0
- average_wait_time: float = 0.0
- class ConnectionPool:
- """
- Async Connection Pool mit automatischer Verwaltung.
- Verwaltet Verbindungen zu einem Host mit:
- - Automatisches Pooling
- - Größenlimits (min/max)
- - Idle-Timeout
- - Lebensdauer-Limits
- - Health-Checks
- Example:
- pool = ConnectionPool(
- host="localhost",
- port=8080,
- config=PoolConfig(min_size=2, max_size=10)
- )
- await pool.start()
- async with pool.acquire() as conn:
- await conn.write(b"Hello")
- response = await conn.read(1024)
- await pool.close()
- """
- def __init__(
- self,
- host: str,
- port: int,
- config: PoolConfig | None = None
- ) -> None:
- """
- Initialisiert den Connection Pool.
- Args:
- host: Zielhost
- port: Zielport
- config: Pool-Konfiguration
- """
- self._host = host
- self._port = port
- self._config = config or PoolConfig()
- self._pool: deque[Connection] = deque()
- self._in_use: set[Connection] = set()
- self._lock = asyncio.Lock()
- self._available = asyncio.Condition(self._lock)
- self._closed = False
- self._maintenance_task: asyncio.Task | None = None
- self._stats = PoolStatistics()
- self._total_wait_time: float = 0.0
- @property
- def host(self) -> str:
- """Zielhost."""
- return self._host
- @property
- def port(self) -> int:
- """Zielport."""
- return self._port
- @property
- def size(self) -> int:
- """Aktuelle Pool-Größe."""
- return len(self._pool) + len(self._in_use)
- @property
- def available(self) -> int:
- """Anzahl verfügbarer Verbindungen."""
- return len(self._pool)
- @property
- def in_use(self) -> int:
- """Anzahl verwendeter Verbindungen."""
- return len(self._in_use)
- async def start(self) -> None:
- """
- Startet den Pool und initialisiert minimale Verbindungen.
- Erstellt die konfigurierten minimalen Verbindungen und
- startet den Maintenance-Task.
- """
- self._closed = False
- # Minimale Verbindungen erstellen
- for _ in range(self._config.min_size):
- conn = await self._create_connection()
- if conn:
- self._pool.append(conn)
- # Maintenance-Task starten
- self._maintenance_task = asyncio.create_task(self._maintenance_loop())
- async def close(self) -> None:
- """
- Schließt den Pool und alle Verbindungen.
- Wartet auf Freigabe verwendeter Verbindungen und schließt
- alle Verbindungen ordnungsgemäß.
- """
- self._closed = True
- # Maintenance-Task stoppen
- if self._maintenance_task:
- self._maintenance_task.cancel()
- try:
- await self._maintenance_task
- except asyncio.CancelledError:
- pass
- # Alle Verbindungen schließen
- async with self._lock:
- for conn in list(self._pool):
- await conn.close()
- self._pool.clear()
- # In-Use-Verbindungen markieren
- for conn in self._in_use:
- await conn.close()
- self._in_use.clear()
- async def acquire(self) -> Connection:
- """
- Holt eine Verbindung aus dem Pool.
- Gibt eine verfügbare Verbindung zurück oder erstellt eine neue,
- wenn das Maximum noch nicht erreicht ist.
- Returns:
- Eine verwendbare Connection
- Raises:
- asyncio.TimeoutError: Wenn kein Timeout innerhalb des Limits verfügbar
- RuntimeError: Wenn Pool geschlossen
- """
- if self._closed:
- raise RuntimeError("Pool ist geschlossen")
- import time
- start_time = time.monotonic()
- try:
- return await asyncio.wait_for(
- self._acquire_impl(),
- timeout=self._config.acquire_timeout
- )
- except asyncio.TimeoutError:
- self._stats.total_timeouts += 1
- raise
- finally:
- wait_time = time.monotonic() - start_time
- self._total_wait_time += wait_time
- self._stats.total_acquires += 1
- if self._stats.total_acquires > 0:
- self._stats.average_wait_time = (
- self._total_wait_time / self._stats.total_acquires
- )
- async def _acquire_impl(self) -> Connection:
- """Interne Acquire-Implementierung."""
- async with self._available:
- while True:
- # Versuche vorhandene Verbindung zu bekommen
- conn = self._get_idle_connection()
- if conn:
- # Validiere wenn konfiguriert
- if self._config.validate_on_acquire:
- if not await self._validate_connection(conn):
- await conn.close()
- continue
- conn.acquire()
- self._in_use.add(conn)
- return conn
- # Keine verfügbare Verbindung - neue erstellen?
- if self.size < self._config.max_size:
- conn = await self._create_connection()
- if conn:
- conn.acquire()
- self._in_use.add(conn)
- return conn
- # Warten auf Freigabe
- await self._available.wait()
- def _get_idle_connection(self) -> Connection | None:
- """Holt eine Idle-Verbindung aus dem Pool."""
- while self._pool:
- conn = self._pool.popleft()
- # Prüfe ob noch gültig
- if conn.info.age > self._config.max_lifetime:
- asyncio.create_task(conn.close())
- continue
- if conn.info.idle_time > self._config.max_idle_time:
- asyncio.create_task(conn.close())
- continue
- if conn.is_usable:
- return conn
- asyncio.create_task(conn.close())
- return None
- async def release(self, conn: Connection) -> None:
- """
- Gibt eine Verbindung zurück in den Pool.
- Args:
- conn: Die freizugebende Verbindung
- """
- async with self._available:
- if conn in self._in_use:
- self._in_use.remove(conn)
- self._stats.total_releases += 1
- # Verbindung prüfen
- if conn.state == ConnectionState.ERROR:
- await conn.close()
- self._stats.total_errors += 1
- elif conn.info.age > self._config.max_lifetime:
- await conn.close()
- elif conn.is_usable:
- conn.release()
- self._pool.append(conn)
- else:
- await conn.close()
- # Warten aufwecken
- self._available.notify()
- async def _create_connection(self) -> Connection | None:
- """Erstellt eine neue Verbindung."""
- conn = Connection(self._host, self._port)
- success = await conn.connect(
- ssl_context=self._config.ssl_context,
- timeout=self._config.connect_timeout
- )
- if success:
- self._stats.total_connections += 1
- return conn
- self._stats.total_errors += 1
- return None
- async def _validate_connection(self, conn: Connection) -> bool:
- """Validiert eine Verbindung."""
- # Einfache Prüfung: Ist die Verbindung noch offen?
- if conn.state != ConnectionState.IDLE:
- return False
- if conn.at_eof():
- return False
- return True
- async def _maintenance_loop(self) -> None:
- """Hintergrund-Task für Pool-Wartung."""
- while not self._closed:
- try:
- await asyncio.sleep(self._config.health_check_interval)
- async with self._lock:
- # Idle-Verbindungen prüfen
- expired = []
- for conn in list(self._pool):
- if conn.info.age > self._config.max_lifetime:
- expired.append(conn)
- elif conn.info.idle_time > self._config.max_idle_time:
- # Behalte mindestens min_size
- if len(self._pool) > self._config.min_size:
- expired.append(conn)
- for conn in expired:
- self._pool.remove(conn)
- await conn.close()
- # Minimale Verbindungen sicherstellen
- while len(self._pool) < self._config.min_size:
- conn = await self._create_connection()
- if conn:
- self._pool.append(conn)
- else:
- break
- # Statistiken aktualisieren
- self._stats.idle_connections = len(self._pool)
- self._stats.in_use_connections = len(self._in_use)
- except asyncio.CancelledError:
- break
- except Exception:
- pass
- def get_statistics(self) -> PoolStatistics:
- """Gibt Pool-Statistiken zurück."""
- stats = PoolStatistics(
- total_connections=self._stats.total_connections,
- idle_connections=len(self._pool),
- in_use_connections=len(self._in_use),
- total_acquires=self._stats.total_acquires,
- total_releases=self._stats.total_releases,
- total_timeouts=self._stats.total_timeouts,
- total_errors=self._stats.total_errors,
- average_wait_time=self._stats.average_wait_time
- )
- return stats
- def to_dict(self) -> dict[str, Any]:
- """Konvertiert zu Dictionary."""
- stats = self.get_statistics()
- return {
- "host": self._host,
- "port": self._port,
- "size": self.size,
- "available": self.available,
- "in_use": self.in_use,
- "config": {
- "min_size": self._config.min_size,
- "max_size": self._config.max_size,
- "max_idle_time": self._config.max_idle_time,
- "max_lifetime": self._config.max_lifetime,
- },
- "statistics": {
- "total_connections": stats.total_connections,
- "total_acquires": stats.total_acquires,
- "total_releases": stats.total_releases,
- "total_timeouts": stats.total_timeouts,
- "total_errors": stats.total_errors,
- "average_wait_time": stats.average_wait_time,
- }
- }
- async def __aenter__(self) -> "ConnectionPool":
- """Context-Manager Entry."""
- await self.start()
- return self
- async def __aexit__(self, exc_type: Any, exc_val: Any, exc_tb: Any) -> None:
- """Context-Manager Exit."""
- await self.close()
- class PooledConnectionContext:
- """
- Context-Manager für gepoolte Verbindungen.
- Ermöglicht einfache Nutzung mit async with.
- """
- def __init__(self, pool: ConnectionPool) -> None:
- self._pool = pool
- self._conn: Connection | None = None
- async def __aenter__(self) -> Connection:
- self._conn = await self._pool.acquire()
- return self._conn
- async def __aexit__(self, exc_type: Any, exc_val: Any, exc_tb: Any) -> None:
- if self._conn:
- await self._pool.release(self._conn)
|