105 lines
3.7 KiB
Python
105 lines
3.7 KiB
Python
"""
|
|
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()
|