""" BrowserPool — N воркеров с очередью визитов. Каждый воркер: 1. Берёт визит из очереди 2. Поднимает свежий браузер (через BrowserService) 3. Делает ровно 1 визит 4. Закрывает браузер 5. Идёт за следующим заданием Пока воркер читает страницу, другие воркеры уже могут запускать браузеры для следующих визитов — стартап перекрывается с чтением. """ import asyncio import logging from dataclasses import dataclass from typing import Optional from services.browser_service import BrowserService, VisitResult logger = logging.getLogger(__name__) @dataclass class _VisitRequest: url: str reading_time: Optional[int] future: asyncio.Future class BrowserPool: """ Пул из N воркеров. Drop-in замена для BrowserService.visit_page(). """ def __init__(self, size: int, browser_service: BrowserService): self._size = size self._bs = browser_service self._queue: asyncio.Queue[Optional[_VisitRequest]] = asyncio.Queue() self._workers: list[asyncio.Task] = [] self._running = False async def start(self) -> None: if self._running: return self._running = True for i in range(self._size): t = asyncio.create_task(self._worker(i), name=f"browser-pool-{i}") self._workers.append(t) logger.info(f"BrowserPool started: {self._size} workers") async def stop(self) -> None: if not self._running: return self._running = False for _ in self._workers: await self._queue.put(None) # sentinel — один на каждого воркера await asyncio.gather(*self._workers, return_exceptions=True) self._workers.clear() logger.info("BrowserPool stopped") async def visit_page(self, url: str, reading_time: int = None) -> VisitResult: """Отправляет визит в очередь и ждёт результата.""" if not self._running: # Пул не запущен — fallback на прямой вызов return await self._bs.visit_page(url, reading_time) loop = asyncio.get_running_loop() future: asyncio.Future[VisitResult] = loop.create_future() await self._queue.put(_VisitRequest(url=url, reading_time=reading_time, future=future)) return await future async def _worker(self, idx: int) -> None: logger.debug(f"Pool worker-{idx} ready") while True: item = await self._queue.get() if item is None: # sentinel → выход break try: result = await self._bs.visit_page(item.url, item.reading_time) if not item.future.done(): item.future.set_result(result) except asyncio.CancelledError: if not item.future.done(): item.future.cancel() break except Exception as e: logger.error(f"Pool worker-{idx} error: {e}", exc_info=True) if not item.future.done(): item.future.set_result( VisitResult(url=item.url, success=False, error=str(e)) ) finally: self._queue.task_done() logger.debug(f"Pool worker-{idx} stopped") @property def size(self) -> int: return self._size @property def queue_size(self) -> int: return self._queue.qsize()