fix payment
This commit is contained in:
@@ -0,0 +1,104 @@
|
||||
"""
|
||||
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()
|
||||
Reference in New Issue
Block a user