pool.py 13 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430
  1. # -*- coding: utf-8 -*-
  2. """
  3. Connection Pool für effiziente Verbindungsverwaltung.
  4. Bietet Pooling mit Größenlimits, Timeouts und Health-Checks.
  5. """
  6. import asyncio
  7. import ssl
  8. from collections import deque
  9. from dataclasses import dataclass
  10. from typing import Any, Callable, Awaitable
  11. from trixy_core.network.pool.connection import Connection, ConnectionState
  12. @dataclass
  13. class PoolConfig:
  14. """
  15. Konfiguration für den Connection Pool.
  16. Attributes:
  17. min_size: Minimale Pool-Größe
  18. max_size: Maximale Pool-Größe
  19. max_idle_time: Maximale Idle-Zeit in Sekunden
  20. max_lifetime: Maximale Lebensdauer einer Verbindung
  21. connect_timeout: Timeout für Verbindungsaufbau
  22. acquire_timeout: Timeout für Pool-Acquire
  23. health_check_interval: Intervall für Health-Checks
  24. validate_on_acquire: Verbindung vor Ausgabe prüfen
  25. ssl_context: Optionaler SSL-Kontext
  26. """
  27. min_size: int = 1
  28. max_size: int = 10
  29. max_idle_time: float = 300.0 # 5 Minuten
  30. max_lifetime: float = 3600.0 # 1 Stunde
  31. connect_timeout: float = 10.0
  32. acquire_timeout: float = 30.0
  33. health_check_interval: float = 30.0
  34. validate_on_acquire: bool = True
  35. ssl_context: ssl.SSLContext | None = None
  36. @dataclass
  37. class PoolStatistics:
  38. """Statistiken des Connection Pools."""
  39. total_connections: int = 0
  40. idle_connections: int = 0
  41. in_use_connections: int = 0
  42. total_acquires: int = 0
  43. total_releases: int = 0
  44. total_timeouts: int = 0
  45. total_errors: int = 0
  46. average_wait_time: float = 0.0
  47. class ConnectionPool:
  48. """
  49. Async Connection Pool mit automatischer Verwaltung.
  50. Verwaltet Verbindungen zu einem Host mit:
  51. - Automatisches Pooling
  52. - Größenlimits (min/max)
  53. - Idle-Timeout
  54. - Lebensdauer-Limits
  55. - Health-Checks
  56. Example:
  57. pool = ConnectionPool(
  58. host="localhost",
  59. port=8080,
  60. config=PoolConfig(min_size=2, max_size=10)
  61. )
  62. await pool.start()
  63. async with pool.acquire() as conn:
  64. await conn.write(b"Hello")
  65. response = await conn.read(1024)
  66. await pool.close()
  67. """
  68. def __init__(
  69. self,
  70. host: str,
  71. port: int,
  72. config: PoolConfig | None = None
  73. ) -> None:
  74. """
  75. Initialisiert den Connection Pool.
  76. Args:
  77. host: Zielhost
  78. port: Zielport
  79. config: Pool-Konfiguration
  80. """
  81. self._host = host
  82. self._port = port
  83. self._config = config or PoolConfig()
  84. self._pool: deque[Connection] = deque()
  85. self._in_use: set[Connection] = set()
  86. self._lock = asyncio.Lock()
  87. self._available = asyncio.Condition(self._lock)
  88. self._closed = False
  89. self._maintenance_task: asyncio.Task | None = None
  90. self._stats = PoolStatistics()
  91. self._total_wait_time: float = 0.0
  92. @property
  93. def host(self) -> str:
  94. """Zielhost."""
  95. return self._host
  96. @property
  97. def port(self) -> int:
  98. """Zielport."""
  99. return self._port
  100. @property
  101. def size(self) -> int:
  102. """Aktuelle Pool-Größe."""
  103. return len(self._pool) + len(self._in_use)
  104. @property
  105. def available(self) -> int:
  106. """Anzahl verfügbarer Verbindungen."""
  107. return len(self._pool)
  108. @property
  109. def in_use(self) -> int:
  110. """Anzahl verwendeter Verbindungen."""
  111. return len(self._in_use)
  112. async def start(self) -> None:
  113. """
  114. Startet den Pool und initialisiert minimale Verbindungen.
  115. Erstellt die konfigurierten minimalen Verbindungen und
  116. startet den Maintenance-Task.
  117. """
  118. self._closed = False
  119. # Minimale Verbindungen erstellen
  120. for _ in range(self._config.min_size):
  121. conn = await self._create_connection()
  122. if conn:
  123. self._pool.append(conn)
  124. # Maintenance-Task starten
  125. self._maintenance_task = asyncio.create_task(self._maintenance_loop())
  126. async def close(self) -> None:
  127. """
  128. Schließt den Pool und alle Verbindungen.
  129. Wartet auf Freigabe verwendeter Verbindungen und schließt
  130. alle Verbindungen ordnungsgemäß.
  131. """
  132. self._closed = True
  133. # Maintenance-Task stoppen
  134. if self._maintenance_task:
  135. self._maintenance_task.cancel()
  136. try:
  137. await self._maintenance_task
  138. except asyncio.CancelledError:
  139. pass
  140. # Alle Verbindungen schließen
  141. async with self._lock:
  142. for conn in list(self._pool):
  143. await conn.close()
  144. self._pool.clear()
  145. # In-Use-Verbindungen markieren
  146. for conn in self._in_use:
  147. await conn.close()
  148. self._in_use.clear()
  149. async def acquire(self) -> Connection:
  150. """
  151. Holt eine Verbindung aus dem Pool.
  152. Gibt eine verfügbare Verbindung zurück oder erstellt eine neue,
  153. wenn das Maximum noch nicht erreicht ist.
  154. Returns:
  155. Eine verwendbare Connection
  156. Raises:
  157. asyncio.TimeoutError: Wenn kein Timeout innerhalb des Limits verfügbar
  158. RuntimeError: Wenn Pool geschlossen
  159. """
  160. if self._closed:
  161. raise RuntimeError("Pool ist geschlossen")
  162. import time
  163. start_time = time.monotonic()
  164. try:
  165. return await asyncio.wait_for(
  166. self._acquire_impl(),
  167. timeout=self._config.acquire_timeout
  168. )
  169. except asyncio.TimeoutError:
  170. self._stats.total_timeouts += 1
  171. raise
  172. finally:
  173. wait_time = time.monotonic() - start_time
  174. self._total_wait_time += wait_time
  175. self._stats.total_acquires += 1
  176. if self._stats.total_acquires > 0:
  177. self._stats.average_wait_time = (
  178. self._total_wait_time / self._stats.total_acquires
  179. )
  180. async def _acquire_impl(self) -> Connection:
  181. """Interne Acquire-Implementierung."""
  182. async with self._available:
  183. while True:
  184. # Versuche vorhandene Verbindung zu bekommen
  185. conn = self._get_idle_connection()
  186. if conn:
  187. # Validiere wenn konfiguriert
  188. if self._config.validate_on_acquire:
  189. if not await self._validate_connection(conn):
  190. await conn.close()
  191. continue
  192. conn.acquire()
  193. self._in_use.add(conn)
  194. return conn
  195. # Keine verfügbare Verbindung - neue erstellen?
  196. if self.size < self._config.max_size:
  197. conn = await self._create_connection()
  198. if conn:
  199. conn.acquire()
  200. self._in_use.add(conn)
  201. return conn
  202. # Warten auf Freigabe
  203. await self._available.wait()
  204. def _get_idle_connection(self) -> Connection | None:
  205. """Holt eine Idle-Verbindung aus dem Pool."""
  206. while self._pool:
  207. conn = self._pool.popleft()
  208. # Prüfe ob noch gültig
  209. if conn.info.age > self._config.max_lifetime:
  210. asyncio.create_task(conn.close())
  211. continue
  212. if conn.info.idle_time > self._config.max_idle_time:
  213. asyncio.create_task(conn.close())
  214. continue
  215. if conn.is_usable:
  216. return conn
  217. asyncio.create_task(conn.close())
  218. return None
  219. async def release(self, conn: Connection) -> None:
  220. """
  221. Gibt eine Verbindung zurück in den Pool.
  222. Args:
  223. conn: Die freizugebende Verbindung
  224. """
  225. async with self._available:
  226. if conn in self._in_use:
  227. self._in_use.remove(conn)
  228. self._stats.total_releases += 1
  229. # Verbindung prüfen
  230. if conn.state == ConnectionState.ERROR:
  231. await conn.close()
  232. self._stats.total_errors += 1
  233. elif conn.info.age > self._config.max_lifetime:
  234. await conn.close()
  235. elif conn.is_usable:
  236. conn.release()
  237. self._pool.append(conn)
  238. else:
  239. await conn.close()
  240. # Warten aufwecken
  241. self._available.notify()
  242. async def _create_connection(self) -> Connection | None:
  243. """Erstellt eine neue Verbindung."""
  244. conn = Connection(self._host, self._port)
  245. success = await conn.connect(
  246. ssl_context=self._config.ssl_context,
  247. timeout=self._config.connect_timeout
  248. )
  249. if success:
  250. self._stats.total_connections += 1
  251. return conn
  252. self._stats.total_errors += 1
  253. return None
  254. async def _validate_connection(self, conn: Connection) -> bool:
  255. """Validiert eine Verbindung."""
  256. # Einfache Prüfung: Ist die Verbindung noch offen?
  257. if conn.state != ConnectionState.IDLE:
  258. return False
  259. if conn.at_eof():
  260. return False
  261. return True
  262. async def _maintenance_loop(self) -> None:
  263. """Hintergrund-Task für Pool-Wartung."""
  264. while not self._closed:
  265. try:
  266. await asyncio.sleep(self._config.health_check_interval)
  267. async with self._lock:
  268. # Idle-Verbindungen prüfen
  269. expired = []
  270. for conn in list(self._pool):
  271. if conn.info.age > self._config.max_lifetime:
  272. expired.append(conn)
  273. elif conn.info.idle_time > self._config.max_idle_time:
  274. # Behalte mindestens min_size
  275. if len(self._pool) > self._config.min_size:
  276. expired.append(conn)
  277. for conn in expired:
  278. self._pool.remove(conn)
  279. await conn.close()
  280. # Minimale Verbindungen sicherstellen
  281. while len(self._pool) < self._config.min_size:
  282. conn = await self._create_connection()
  283. if conn:
  284. self._pool.append(conn)
  285. else:
  286. break
  287. # Statistiken aktualisieren
  288. self._stats.idle_connections = len(self._pool)
  289. self._stats.in_use_connections = len(self._in_use)
  290. except asyncio.CancelledError:
  291. break
  292. except Exception:
  293. pass
  294. def get_statistics(self) -> PoolStatistics:
  295. """Gibt Pool-Statistiken zurück."""
  296. stats = PoolStatistics(
  297. total_connections=self._stats.total_connections,
  298. idle_connections=len(self._pool),
  299. in_use_connections=len(self._in_use),
  300. total_acquires=self._stats.total_acquires,
  301. total_releases=self._stats.total_releases,
  302. total_timeouts=self._stats.total_timeouts,
  303. total_errors=self._stats.total_errors,
  304. average_wait_time=self._stats.average_wait_time
  305. )
  306. return stats
  307. def to_dict(self) -> dict[str, Any]:
  308. """Konvertiert zu Dictionary."""
  309. stats = self.get_statistics()
  310. return {
  311. "host": self._host,
  312. "port": self._port,
  313. "size": self.size,
  314. "available": self.available,
  315. "in_use": self.in_use,
  316. "config": {
  317. "min_size": self._config.min_size,
  318. "max_size": self._config.max_size,
  319. "max_idle_time": self._config.max_idle_time,
  320. "max_lifetime": self._config.max_lifetime,
  321. },
  322. "statistics": {
  323. "total_connections": stats.total_connections,
  324. "total_acquires": stats.total_acquires,
  325. "total_releases": stats.total_releases,
  326. "total_timeouts": stats.total_timeouts,
  327. "total_errors": stats.total_errors,
  328. "average_wait_time": stats.average_wait_time,
  329. }
  330. }
  331. async def __aenter__(self) -> "ConnectionPool":
  332. """Context-Manager Entry."""
  333. await self.start()
  334. return self
  335. async def __aexit__(self, exc_type: Any, exc_val: Any, exc_tb: Any) -> None:
  336. """Context-Manager Exit."""
  337. await self.close()
  338. class PooledConnectionContext:
  339. """
  340. Context-Manager für gepoolte Verbindungen.
  341. Ermöglicht einfache Nutzung mit async with.
  342. """
  343. def __init__(self, pool: ConnectionPool) -> None:
  344. self._pool = pool
  345. self._conn: Connection | None = None
  346. async def __aenter__(self) -> Connection:
  347. self._conn = await self._pool.acquire()
  348. return self._conn
  349. async def __aexit__(self, exc_type: Any, exc_val: Any, exc_tb: Any) -> None:
  350. if self._conn:
  351. await self._pool.release(self._conn)