From ef8c5bc7d6d85820b791cb099fc9bf71681d8187 Mon Sep 17 00:00:00 2001 From: Karim shoair Date: Mon, 2 Feb 2026 00:16:22 +0200 Subject: [PATCH] feat(spiders/fetchers): Adding proxy rotation logic and change retry logic - User passes a single proxy to the browser session, and it will keep using the same context. Pass a proxy manager that automatically refreshes the IP with each proxy, and you get speed but sacrifice some stealth. - User imports the proxy rotator class and passes proxies to it and a rotation strategy (Round robin by default). Then pass the instance to the session class, which will create a context and a tab for each proxy returned by the rotator. This way, you sacrifice speed a bit, but you get maximum stealth since each context is created with the proxy that will be used. - Normal requests use what you pass without issues, of course. - All errors are retried now, and proxies are rotated on retries automatically. --- scrapling/core/_types.py | 2 + scrapling/engines/_browsers/_base.py | 191 +++++++++++-- scrapling/engines/_browsers/_controllers.py | 235 ++++++++------- scrapling/engines/_browsers/_stealth.py | 283 ++++++++++--------- scrapling/engines/_browsers/_types.py | 3 + scrapling/engines/_browsers/_validators.py | 7 + scrapling/engines/constants.py | 2 - scrapling/engines/static.py | 75 ++++- scrapling/engines/toolbelt/__init__.py | 2 + scrapling/engines/toolbelt/proxy_rotation.py | 104 +++++++ 10 files changed, 623 insertions(+), 281 deletions(-) create mode 100644 scrapling/engines/toolbelt/proxy_rotation.py diff --git a/scrapling/core/_types.py b/scrapling/core/_types.py index 48a2eec..f19e706 100644 --- a/scrapling/core/_types.py +++ b/scrapling/core/_types.py @@ -33,6 +33,8 @@ from typing import ( SupportsIndex, ) +# Proxy can be a string URL or a dict (Playwright format: {"server": "...", "username": "...", "password": "..."}) +ProxyType = Union[str, Dict[str, str]] SUPPORTED_HTTP_METHODS = Literal["GET", "POST", "PUT", "DELETE"] SelectorWaitStates = Literal["attached", "detached", "hidden", "visible"] PageLoadStates = Literal["commit", "domcontentloaded", "load", "networkidle"] diff --git a/scrapling/engines/_browsers/_base.py b/scrapling/engines/_browsers/_base.py index 810c51c..5936652 100644 --- a/scrapling/engines/_browsers/_base.py +++ b/scrapling/engines/_browsers/_base.py @@ -1,9 +1,11 @@ from time import time from asyncio import sleep as asyncio_sleep, Lock +from contextlib import contextmanager, asynccontextmanager from playwright.sync_api._generated import Page from playwright.sync_api import ( Frame, + Browser, BrowserContext, Playwright, Response as SyncPlaywrightResponse, @@ -11,18 +13,31 @@ from playwright.sync_api import ( from playwright.async_api._generated import Page as AsyncPage from playwright.async_api import ( Frame as AsyncFrame, + Browser as AsyncBrowser, Playwright as AsyncPlaywright, Response as AsyncPlaywrightResponse, BrowserContext as AsyncBrowserContext, ) from playwright._impl._errors import Error as PlaywrightError -from ._page import PageInfo, PagePool from scrapling.parser import Selector -from ._validators import validate, PlaywrightConfig, StealthConfig -from ._config_tools import __default_chrome_useragent__, __default_useragent__ -from scrapling.engines.toolbelt.navigation import intercept_route, async_intercept_route -from scrapling.core._types import Any, cast, Dict, List, Optional, Callable, TYPE_CHECKING, overload, Tuple +from scrapling.engines._browsers._page import PageInfo, PagePool +from scrapling.engines._browsers._validators import validate, PlaywrightConfig, StealthConfig +from scrapling.engines._browsers._config_tools import __default_chrome_useragent__, __default_useragent__ +from scrapling.engines.toolbelt.navigation import construct_proxy_dict, intercept_route, async_intercept_route +from scrapling.core._types import ( + Any, + Dict, + List, + Optional, + Callable, + TYPE_CHECKING, + overload, + Tuple, + ProxyType, + Generator, + AsyncGenerator, +) from scrapling.engines.constants import ( DEFAULT_STEALTH_FLAGS, HARMFUL_DEFAULT_ARGS, @@ -31,12 +46,19 @@ from scrapling.engines.constants import ( class SyncSession: + _config: "PlaywrightConfig | StealthConfig" + _context_options: Dict[str, Any] + + def _build_context_with_proxy(self, proxy: Optional[ProxyType] = None) -> Dict[str, Any]: + raise NotImplementedError # pragma: no cover + def __init__(self, max_pages: int = 1): self.max_pages = max_pages self.page_pool = PagePool(max_pages) self._max_wait_for_page = 60 self.playwright: Playwright | Any = None self.context: BrowserContext | Any = None + self.browser: Optional[Browser] = None self._is_alive = False def start(self): @@ -51,6 +73,10 @@ class SyncSession: self.context.close() self.context = None + if self.browser: + self.browser.close() + self.browser = None + if self.playwright: self.playwright.stop() self.playwright = None # pyright: ignore @@ -64,17 +90,28 @@ class SyncSession: def __exit__(self, exc_type, exc_val, exc_tb): self.close() + def _initialize_context(self, config: PlaywrightConfig | StealthConfig, ctx: BrowserContext) -> BrowserContext: + """Initialize the browser context.""" + if config.init_script: + ctx.add_init_script(path=config.init_script) + + if config.cookies: # pragma: no cover + ctx.add_cookies(config.cookies) + + return ctx + def _get_page( self, timeout: int | float, extra_headers: Optional[Dict[str, str]], disable_resources: bool, + context: Optional[BrowserContext] = None, ) -> PageInfo[Page]: # pragma: no cover """Get a new page to use""" - # No need to check if a page is available or not in sync code because the code blocked before reaching here till the page closed, ofc. - assert self.context is not None, "Browser context not initialized" - page = self.context.new_page() + ctx = context if context is not None else self.context + assert ctx is not None, "Browser context not initialized" + page = ctx.new_page() page.set_default_navigation_timeout(timeout) page.set_default_timeout(timeout) if extra_headers: @@ -129,14 +166,52 @@ class SyncSession: return handle_response + @contextmanager + def _page_generator( + self, + timeout: int | float, + extra_headers: Optional[Dict[str, str]], + disable_resources: bool, + proxy: Optional[ProxyType] = None, + ) -> Generator["PageInfo[Page]", None, None]: + """Acquire a page - either from persistent context or fresh context with proxy.""" + if self._config.proxy_rotator: + # Rotation mode: create fresh context with the provided proxy + if not self.browser: # pragma: no cover + raise RuntimeError("Browser not initialized for proxy rotation mode") + context_options = self._build_context_with_proxy(proxy) + context: BrowserContext = self.browser.new_context(**context_options) + + try: + context = self._initialize_context(self._config, context) + page_info = self._get_page(timeout, extra_headers, disable_resources, context=context) + yield page_info + finally: + context.close() + else: + # Standard mode: use PagePool with persistent context + page_info = self._get_page(timeout, extra_headers, disable_resources) + try: + yield page_info + finally: + page_info.page.close() + self.page_pool.pages.remove(page_info) + class AsyncSession: + _config: "PlaywrightConfig | StealthConfig" + _context_options: Dict[str, Any] + + def _build_context_with_proxy(self, proxy: Optional[ProxyType] = None) -> Dict[str, Any]: + raise NotImplementedError # pragma: no cover + def __init__(self, max_pages: int = 1): self.max_pages = max_pages self.page_pool = PagePool(max_pages) self._max_wait_for_page = 60 self.playwright: AsyncPlaywright | Any = None self.context: AsyncBrowserContext | Any = None + self.browser: Optional[AsyncBrowser] = None self._is_alive = False self._lock = Lock() @@ -152,6 +227,10 @@ class AsyncSession: await self.context.close() self.context = None # pyright: ignore + if self.browser: + await self.browser.close() + self.browser = None + if self.playwright: await self.playwright.stop() self.playwright = None # pyright: ignore @@ -165,19 +244,34 @@ class AsyncSession: async def __aexit__(self, exc_type, exc_val, exc_tb): await self.close() + async def _initialize_context( + self, config: PlaywrightConfig | StealthConfig, ctx: AsyncBrowserContext + ) -> AsyncBrowserContext: + """Initialize the browser context.""" + if config.init_script: # pragma: no cover + await ctx.add_init_script(path=config.init_script) + + if config.cookies: # pragma: no cover + await ctx.add_cookies(config.cookies) + + return ctx + async def _get_page( self, timeout: int | float, extra_headers: Optional[Dict[str, str]], disable_resources: bool, + context: Optional[AsyncBrowserContext] = None, ) -> PageInfo[AsyncPage]: # pragma: no cover """Get a new page to use""" + ctx = context if context is not None else self.context if TYPE_CHECKING: - assert self.context is not None, "Browser context not initialized" + assert ctx is not None, "Browser context not initialized" async with self._lock: # If we're at max capacity after cleanup, wait for busy pages to finish - if self.page_pool.pages_count >= self.max_pages: + if context is None and self.page_pool.pages_count >= self.max_pages: + # Only applies when using persistent context start_time = time() while time() - start_time < self._max_wait_for_page: await asyncio_sleep(0.05) @@ -188,7 +282,7 @@ class AsyncSession: f"No pages finished to clear place in the pool within the {self._max_wait_for_page}s timeout period" ) - page = await self.context.new_page() + page = await ctx.new_page() page.set_default_navigation_timeout(timeout) page.set_default_timeout(timeout) if extra_headers: @@ -241,6 +335,37 @@ class AsyncSession: return handle_response + @asynccontextmanager + async def _page_generator( + self, + timeout: int | float, + extra_headers: Optional[Dict[str, str]], + disable_resources: bool, + proxy: Optional[ProxyType] = None, + ) -> AsyncGenerator["PageInfo[AsyncPage]", None]: + """Acquire a page - either from persistent context or fresh context with proxy.""" + if self._config.proxy_rotator: + # Rotation mode: create fresh context with the provided proxy + if not self.browser: # pragma: no cover + raise RuntimeError("Browser not initialized for proxy rotation mode") + context_options = self._build_context_with_proxy(proxy) + context: AsyncBrowserContext = await self.browser.new_context(**context_options) + + try: + context = await self._initialize_context(self._config, context) + page_info = await self._get_page(timeout, extra_headers, disable_resources, context=context) + yield page_info + finally: + await context.close() + else: + # Standard mode: use PagePool with persistent context + page_info = await self._get_page(timeout, extra_headers, disable_resources) + try: + yield page_info + finally: + await page_info.page.close() + self.page_pool.pages.remove(page_info) + class BaseSessionMixin: @overload @@ -254,7 +379,7 @@ class BaseSessionMixin: ) -> PlaywrightConfig | StealthConfig: # Dark color scheme bypasses the 'prefersLightColor' check in creepjs self._context_options: Dict[str, Any] = {"color_scheme": "dark", "device_scale_factor": 2} - self._launch_options: Dict[str, Any] = self._context_options | { + self._browser_options: Dict[str, Any] = { "args": DEFAULT_FLAGS, "ignore_default_args": HARMFUL_DEFAULT_ARGS, } @@ -269,7 +394,7 @@ class BaseSessionMixin: return config def __generate_options__(self, extra_flags: Tuple | None = None) -> None: - config = cast(PlaywrightConfig, getattr(self, "_config", None)) + config: PlaywrightConfig | StealthConfig = self._config # type: ignore[has-type] self._context_options.update( { "proxy": config.proxy, @@ -287,36 +412,40 @@ class BaseSessionMixin: ) if not config.cdp_url: - self._launch_options |= self._context_options - self._context_options = {} - flags = self._launch_options["args"] + flags = self._browser_options["args"] if config.extra_flags or extra_flags: flags = list(set(flags + (config.extra_flags or extra_flags))) - self._launch_options.update( + self._browser_options.update( { "args": flags, "headless": config.headless, - "user_data_dir": config.user_data_dir, "channel": "chrome" if config.real_chrome else "chromium", } ) - if config.additional_args: - self._launch_options.update(config.additional_args) + self._user_data_dir = config.user_data_dir else: - # while `context_options` is left to be used when cdp mode is enabled - self._launch_options = dict() - if config.additional_args: - self._context_options.update(config.additional_args) + self._browser_options = {} - @staticmethod - def _is_retriable(error: Exception) -> bool: - """Check if an error is retriable (transient network/timeout issues).""" - if isinstance(error, TimeoutError): - return True - error_msg = str(error).lower() - return "net::" in error_msg or "failed to get response" in error_msg + if config.additional_args: + self._context_options.update(config.additional_args) + + def _build_context_with_proxy(self, proxy: Optional[ProxyType] = None) -> Dict[str, Any]: + """ + Build context options with a specific proxy for rotation mode. + + :param proxy: Proxy URL string or Playwright-style proxy dict to use for this context. + :return: Dictionary of context options for browser.new_context(). + """ + + context_options = self._context_options.copy() + + # Override proxy if provided + if proxy: + context_options["proxy"] = construct_proxy_dict(proxy) + + return context_options class DynamicSessionMixin(BaseSessionMixin): diff --git a/scrapling/engines/_browsers/_controllers.py b/scrapling/engines/_browsers/_controllers.py index a9dbe85..34a0b25 100644 --- a/scrapling/engines/_browsers/_controllers.py +++ b/scrapling/engines/_browsers/_controllers.py @@ -12,12 +12,13 @@ from playwright.async_api import ( ) from scrapling.core.utils import log -from scrapling.core._types import Unpack, TYPE_CHECKING -from ._types import PlaywrightSession, PlaywrightFetchParams -from ._base import SyncSession, AsyncSession, DynamicSessionMixin -from ._validators import validate_fetch as _validate, PlaywrightConfig +from scrapling.core._types import Unpack +from scrapling.engines.toolbelt.proxy_rotation import is_proxy_error from scrapling.engines.toolbelt.convertor import Response, ResponseFactory from scrapling.engines.toolbelt.fingerprints import generate_convincing_referer +from scrapling.engines._browsers._types import PlaywrightSession, PlaywrightFetchParams +from scrapling.engines._browsers._base import SyncSession, AsyncSession, DynamicSessionMixin +from scrapling.engines._browsers._validators import validate_fetch as _validate, PlaywrightConfig class DynamicSession(SyncSession, DynamicSessionMixin): @@ -26,7 +27,9 @@ class DynamicSession(SyncSession, DynamicSessionMixin): __slots__ = ( "_config", "_context_options", - "_launch_options", + "_browser_options", + "_user_data_dir", + "_headers_keys", "max_pages", "page_pool", "_max_wait_for_page", @@ -73,16 +76,19 @@ class DynamicSession(SyncSession, DynamicSessionMixin): try: if self._config.cdp_url: # pragma: no cover - browser = self.playwright.chromium.connect_over_cdp(endpoint_url=self._config.cdp_url) - self.context = browser.new_context(**self._context_options) + self.browser = self.playwright.chromium.connect_over_cdp(endpoint_url=self._config.cdp_url) + if not self._config.proxy_rotator and self.browser: + self.context = self.browser.new_context(**self._context_options) + elif self._config.proxy_rotator: + self.browser = self.playwright.chromium.launch(**self._browser_options) else: - self.context = self.playwright.chromium.launch_persistent_context(**self._launch_options) + persistent_options = ( + self._browser_options | self._context_options | {"user_data_dir": self._user_data_dir} + ) + self.context = self.playwright.chromium.launch_persistent_context(**persistent_options) - if self._config.init_script: # pragma: no cover - self.context.add_init_script(path=self._config.init_script) - - if self._config.cookies: # pragma: no cover - self.context.add_cookies(self._config.cookies) + if self.context: + self.context = self._initialize_context(self._config, self.context) self._is_alive = True except Exception: @@ -123,60 +129,73 @@ class DynamicSession(SyncSession, DynamicSessionMixin): ) for attempt in range(self._config.retries): - page_info = self._get_page(params.timeout, params.extra_headers, params.disable_resources) - final_response = [None] - handle_response = self._create_response_handler(page_info, final_response) + proxy = self._config.proxy_rotator.get_proxy() if self._config.proxy_rotator else None - try: # pragma: no cover - page_info.page.on("response", handle_response) - first_response = page_info.page.goto(url, referer=referer) - self._wait_for_page_stability(page_info.page, params.load_dom, params.network_idle) + with self._page_generator( + params.timeout, params.extra_headers, params.disable_resources, proxy + ) as page_info: + final_response = [None] + page = page_info.page + page.on("response", self._create_response_handler(page_info, final_response)) - if not first_response: - raise RuntimeError(f"Failed to get response for {url}") + try: + first_response = page.goto(url, referer=referer) + self._wait_for_page_stability(page, params.load_dom, params.network_idle) - if params.page_action: - try: - _ = params.page_action(page_info.page) - except Exception as e: # pragma: no cover - log.error(f"Error executing page_action: {e}") + if not first_response: + raise RuntimeError(f"Failed to get response for {url}") - if params.wait_selector: - try: - waiter: Locator = page_info.page.locator(params.wait_selector) - waiter.first.wait_for(state=params.wait_selector_state) - self._wait_for_page_stability(page_info.page, params.load_dom, params.network_idle) - except Exception as e: # pragma: no cover - log.error(f"Error waiting for selector {params.wait_selector}: {e}") + if params.page_action: + try: + _ = params.page_action(page) + except Exception as e: # pragma: no cover + log.error(f"Error executing page_action: {e}") - page_info.page.wait_for_timeout(params.wait) + if params.wait_selector: + try: + waiter: Locator = page.locator(params.wait_selector) + waiter.first.wait_for(state=params.wait_selector_state) + self._wait_for_page_stability(page, params.load_dom, params.network_idle) + except Exception as e: # pragma: no cover + log.error(f"Error waiting for selector {params.wait_selector}: {e}") - response = ResponseFactory.from_playwright_response( - page_info.page, first_response, final_response[0], params.selector_config - ) + page.wait_for_timeout(params.wait) - page_info.page.close() - self.page_pool.pages.remove(page_info) - return response + response = ResponseFactory.from_playwright_response( + page, first_response, final_response[0], params.selector_config + ) + return response - except Exception as e: - page_info.mark_error() - page_info.page.close() - self.page_pool.pages.remove(page_info) + except Exception as e: + page_info.mark_error() + if attempt < self._config.retries - 1: + if is_proxy_error(e): + log.warning( + f"Proxy '{proxy}' failed (attempt {attempt + 1}) | Retrying in {self._config.retry_delay}s..." + ) + else: + log.warning( + f"Attempt {attempt + 1} failed: {e}. Retrying in {self._config.retry_delay}s..." + ) + time_sleep(self._config.retry_delay) + else: + log.error(f"Failed after {self._config.retries} attempts: {e}") + raise - if attempt < self._config.retries - 1 and self._is_retriable(e): - log.warning(f"Attempt {attempt + 1} failed: {e}. Retrying in {self._config.retry_delay}s...") - time_sleep(self._config.retry_delay) - else: - raise - - # For type checking purposes only - raise AssertionError("Unreachable: retry loop must return or raise") # pragma: no cover + raise RuntimeError("Request failed") # pragma: no cover class AsyncDynamicSession(AsyncSession, DynamicSessionMixin): """An async Browser session manager with page pooling, it's using a persistent browser Context by default with a temporary user profile directory.""" + __slots__ = ( + "_config", + "_context_options", + "_browser_options", + "_user_data_dir", + "_headers_keys", + ) + def __init__(self, **kwargs: Unpack[PlaywrightSession]): """A Browser session manager with page pooling @@ -216,18 +235,21 @@ class AsyncDynamicSession(AsyncSession, DynamicSessionMixin): self.playwright = await async_playwright().start() try: if self._config.cdp_url: - browser = await self.playwright.chromium.connect_over_cdp(endpoint_url=self._config.cdp_url) - self.context: AsyncBrowserContext = await browser.new_context(**self._context_options) + self.browser = await self.playwright.chromium.connect_over_cdp(endpoint_url=self._config.cdp_url) + if not self._config.proxy_rotator and self.browser: + self.context: AsyncBrowserContext = await self.browser.new_context(**self._context_options) + elif self._config.proxy_rotator: + self.browser = await self.playwright.chromium.launch(**self._browser_options) else: + persistent_options = ( + self._browser_options | self._context_options | {"user_data_dir": self._user_data_dir} + ) self.context: AsyncBrowserContext = await self.playwright.chromium.launch_persistent_context( - **self._launch_options + **persistent_options ) - if self._config.init_script: # pragma: no cover - await self.context.add_init_script(path=self._config.init_script) - - if self._config.cookies: - await self.context.add_cookies(self._config.cookies) # pyright: ignore + if self.context: + self.context = await self._initialize_context(self._config, self.context) self._is_alive = True except Exception: @@ -269,58 +291,57 @@ class AsyncDynamicSession(AsyncSession, DynamicSessionMixin): ) for attempt in range(self._config.retries): - page_info = await self._get_page(params.timeout, params.extra_headers, params.disable_resources) - final_response = [None] - handle_response = self._create_response_handler(page_info, final_response) + proxy = self._config.proxy_rotator.get_proxy() if self._config.proxy_rotator else None - if TYPE_CHECKING: - from playwright.async_api import Page as async_Page + async with self._page_generator( + params.timeout, params.extra_headers, params.disable_resources, proxy + ) as page_info: + final_response = [None] + page = page_info.page + page.on("response", self._create_response_handler(page_info, final_response)) - if not isinstance(page_info.page, async_Page): - raise TypeError + try: + first_response = await page.goto(url, referer=referer) + await self._wait_for_page_stability(page, params.load_dom, params.network_idle) - try: - page_info.page.on("response", handle_response) - first_response = await page_info.page.goto(url, referer=referer) - await self._wait_for_page_stability(page_info.page, params.load_dom, params.network_idle) + if not first_response: + raise RuntimeError(f"Failed to get response for {url}") - if not first_response: - raise RuntimeError(f"Failed to get response for {url}") + if params.page_action: + try: + _ = await params.page_action(page) + except Exception as e: # pragma: no cover + log.error(f"Error executing page_action: {e}") - if params.page_action: - try: - _ = await params.page_action(page_info.page) - except Exception as e: - log.error(f"Error executing page_action: {e}") + if params.wait_selector: + try: + waiter: AsyncLocator = page.locator(params.wait_selector) + await waiter.first.wait_for(state=params.wait_selector_state) + await self._wait_for_page_stability(page, params.load_dom, params.network_idle) + except Exception as e: # pragma: no cover + log.error(f"Error waiting for selector {params.wait_selector}: {e}") - if params.wait_selector: - try: - waiter: AsyncLocator = page_info.page.locator(params.wait_selector) - await waiter.first.wait_for(state=params.wait_selector_state) - await self._wait_for_page_stability(page_info.page, params.load_dom, params.network_idle) - except Exception as e: - log.error(f"Error waiting for selector {params.wait_selector}: {e}") + await page.wait_for_timeout(params.wait) - await page_info.page.wait_for_timeout(params.wait) + response = await ResponseFactory.from_async_playwright_response( + page, first_response, final_response[0], params.selector_config + ) + return response - response = await ResponseFactory.from_async_playwright_response( - page_info.page, first_response, final_response[0], params.selector_config - ) + except Exception as e: + page_info.mark_error() + if attempt < self._config.retries - 1: + if is_proxy_error(e): + log.warning( + f"Proxy '{proxy}' failed (attempt {attempt + 1}) | Retrying in {self._config.retry_delay}s..." + ) + else: + log.warning( + f"Attempt {attempt + 1} failed: {e}. Retrying in {self._config.retry_delay}s..." + ) + await asyncio_sleep(self._config.retry_delay) + else: + log.error(f"Failed after {self._config.retries} attempts: {e}") + raise - await page_info.page.close() - self.page_pool.pages.remove(page_info) - return response - - except Exception as e: # pragma: no cover - page_info.mark_error() - await page_info.page.close() - self.page_pool.pages.remove(page_info) - - if attempt < self._config.retries - 1 and self._is_retriable(e): - log.warning(f"Attempt {attempt + 1} failed: {e}. Retrying in {self._config.retry_delay}s...") - await asyncio_sleep(self._config.retry_delay) - else: - raise - - # For type checking purposes only - raise AssertionError("Unreachable: retry loop must return or raise") # pragma: no cover + raise RuntimeError("Request failed") # pragma: no cover diff --git a/scrapling/engines/_browsers/_stealth.py b/scrapling/engines/_browsers/_stealth.py index 62c6455..fc76d26 100644 --- a/scrapling/engines/_browsers/_stealth.py +++ b/scrapling/engines/_browsers/_stealth.py @@ -3,10 +3,7 @@ from re import compile as re_compile from time import sleep as time_sleep from asyncio import sleep as asyncio_sleep -from playwright.sync_api import ( - Locator, - Page, -) +from playwright.sync_api import Locator, Page, BrowserContext from playwright.async_api import ( Page as async_Page, Locator as AsyncLocator, @@ -17,12 +14,13 @@ from patchright.async_api import async_playwright from scrapling.core.utils import log from scrapling.core._types import Any, Unpack -from ._config_tools import _compiled_stealth_scripts -from ._types import StealthSession, StealthFetchParams -from ._base import SyncSession, AsyncSession, StealthySessionMixin -from ._validators import validate_fetch as _validate, StealthConfig +from scrapling.engines.toolbelt.proxy_rotation import is_proxy_error from scrapling.engines.toolbelt.convertor import Response, ResponseFactory from scrapling.engines.toolbelt.fingerprints import generate_convincing_referer +from scrapling.engines._browsers._config_tools import _compiled_stealth_scripts +from scrapling.engines._browsers._types import StealthSession, StealthFetchParams +from scrapling.engines._browsers._base import SyncSession, AsyncSession, StealthySessionMixin +from scrapling.engines._browsers._validators import validate_fetch as _validate, StealthConfig __CF_PATTERN__ = re_compile("challenges.cloudflare.com/cdn-cgi/challenge-platform/.*") @@ -33,7 +31,9 @@ class StealthySession(SyncSession, StealthySessionMixin): __slots__ = ( "_config", "_context_options", - "_launch_options", + "_browser_options", + "_user_data_dir", + "_headers_keys", "max_pages", "page_pool", "_max_wait_for_page", @@ -84,19 +84,20 @@ class StealthySession(SyncSession, StealthySessionMixin): try: if self._config.cdp_url: # pragma: no cover - browser = self.playwright.chromium.connect_over_cdp(endpoint_url=self._config.cdp_url) - self.context = browser.new_context(**self._context_options) + self.browser = self.playwright.chromium.connect_over_cdp(endpoint_url=self._config.cdp_url) + if not self._config.proxy_rotator: + assert self.browser is not None + self.context = self.browser.new_context(**self._context_options) + elif self._config.proxy_rotator: + self.browser = self.playwright.chromium.launch(**self._browser_options) else: - self.context = self.playwright.chromium.launch_persistent_context(**self._launch_options) + persistent_options = ( + self._browser_options | self._context_options | {"user_data_dir": self._user_data_dir} + ) + self.context = self.playwright.chromium.launch_persistent_context(**persistent_options) - for script in _compiled_stealth_scripts(): - self.context.add_init_script(script=script) - - if self._config.init_script: # pragma: no cover - self.context.add_init_script(path=self._config.init_script) - - if self._config.cookies: # pragma: no cover - self.context.add_cookies(self._config.cookies) + if self.context: + self.context = self._initialize_context(self._config, self.context) self._is_alive = True except Exception: @@ -107,6 +108,14 @@ class StealthySession(SyncSession, StealthySessionMixin): else: raise RuntimeError("Session has been already started") + def _initialize_context(self, config, ctx: BrowserContext) -> BrowserContext: + """Initialize the browser context.""" + for script in _compiled_stealth_scripts(): + ctx.add_init_script(script=script) + + ctx = super()._initialize_context(config, ctx) + return ctx + def _cloudflare_solver(self, page: Page) -> None: # pragma: no cover """Solve the cloudflare challenge displayed on the playwright page passed @@ -209,70 +218,78 @@ class StealthySession(SyncSession, StealthySessionMixin): ) for attempt in range(self._config.retries): - page_info = self._get_page(params.timeout, params.extra_headers, params.disable_resources) - final_response = [None] - handle_response = self._create_response_handler(page_info, final_response) + proxy = self._config.proxy_rotator.get_proxy() if self._config.proxy_rotator else None - try: # pragma: no cover - # Navigate to URL and wait for a specified state - page_info.page.on("response", handle_response) - first_response = page_info.page.goto(url, referer=referer) - self._wait_for_page_stability(page_info.page, params.load_dom, params.network_idle) + with self._page_generator( + params.timeout, params.extra_headers, params.disable_resources, proxy + ) as page_info: + final_response = [None] + page = page_info.page + page.on("response", self._create_response_handler(page_info, final_response)) - if not first_response: - raise RuntimeError(f"Failed to get response for {url}") + try: + first_response = page.goto(url, referer=referer) + self._wait_for_page_stability(page, params.load_dom, params.network_idle) - if params.solve_cloudflare: - self._cloudflare_solver(page_info.page) - # Make sure the page is fully loaded after the captcha - self._wait_for_page_stability(page_info.page, params.load_dom, params.network_idle) + if not first_response: + raise RuntimeError(f"Failed to get response for {url}") - if params.page_action: - try: - _ = params.page_action(page_info.page) - except Exception as e: # pragma: no cover - log.error(f"Error executing page_action: {e}") + if params.solve_cloudflare: + self._cloudflare_solver(page) + # Make sure the page is fully loaded after the captcha + self._wait_for_page_stability(page, params.load_dom, params.network_idle) - if params.wait_selector: - try: - waiter: Locator = page_info.page.locator(params.wait_selector) - waiter.first.wait_for(state=params.wait_selector_state) - # Wait again after waiting for the selector, helpful with protections like Cloudflare - self._wait_for_page_stability(page_info.page, params.load_dom, params.network_idle) - except Exception as e: # pragma: no cover - log.error(f"Error waiting for selector {params.wait_selector}: {e}") + if params.page_action: + try: + _ = params.page_action(page) + except Exception as e: # pragma: no cover + log.error(f"Error executing page_action: {e}") - page_info.page.wait_for_timeout(params.wait) + if params.wait_selector: + try: + waiter: Locator = page.locator(params.wait_selector) + waiter.first.wait_for(state=params.wait_selector_state) + self._wait_for_page_stability(page, params.load_dom, params.network_idle) + except Exception as e: # pragma: no cover + log.error(f"Error waiting for selector {params.wait_selector}: {e}") - # Create response object - response = ResponseFactory.from_playwright_response( - page_info.page, first_response, final_response[0], params.selector_config - ) + page.wait_for_timeout(params.wait) - # Close the page to free up resources - page_info.page.close() - self.page_pool.pages.remove(page_info) + response = ResponseFactory.from_playwright_response( + page, first_response, final_response[0], params.selector_config + ) + return response - return response + except Exception as e: + page_info.mark_error() + if attempt < self._config.retries - 1: + if is_proxy_error(e): + log.warning( + f"Proxy '{proxy}' failed (attempt {attempt + 1}) | Retrying in {self._config.retry_delay}s..." + ) + else: + log.warning( + f"Attempt {attempt + 1} failed: {e}. Retrying in {self._config.retry_delay}s..." + ) + time_sleep(self._config.retry_delay) + else: + log.error(f"Failed after {self._config.retries} attempts: {e}") + raise - except Exception as e: - page_info.mark_error() - page_info.page.close() - self.page_pool.pages.remove(page_info) - - if attempt < self._config.retries - 1 and self._is_retriable(e): - log.warning(f"Attempt {attempt + 1} failed: {e}. Retrying in {self._config.retry_delay}s...") - time_sleep(self._config.retry_delay) - else: - raise - - # For type checking purposes only - raise AssertionError("Unreachable: retry loop must return or raise") # pragma: no cover + raise RuntimeError("Request failed") # pragma: no cover class AsyncStealthySession(AsyncSession, StealthySessionMixin): """An async Stealthy Browser session manager with page pooling.""" + __slots__ = ( + "_config", + "_context_options", + "_browser_options", + "_user_data_dir", + "_headers_keys", + ) + def __init__(self, **kwargs: Unpack[StealthSession]): """A Browser session manager with page pooling, it's using a persistent browser Context by default with a temporary user profile directory. @@ -315,21 +332,22 @@ class AsyncStealthySession(AsyncSession, StealthySessionMixin): self.playwright = await async_playwright().start() try: if self._config.cdp_url: - browser = await self.playwright.chromium.connect_over_cdp(endpoint_url=self._config.cdp_url) - self.context: AsyncBrowserContext = await browser.new_context(**self._context_options) + self.browser = await self.playwright.chromium.connect_over_cdp(endpoint_url=self._config.cdp_url) + if not self._config.proxy_rotator: + assert self.browser is not None + self.context: AsyncBrowserContext = await self.browser.new_context(**self._context_options) + elif self._config.proxy_rotator: + self.browser = await self.playwright.chromium.launch(**self._browser_options) else: + persistent_options = ( + self._browser_options | self._context_options | {"user_data_dir": self._user_data_dir} + ) self.context: AsyncBrowserContext = await self.playwright.chromium.launch_persistent_context( - **self._launch_options + **persistent_options ) - for script in _compiled_stealth_scripts(): - await self.context.add_init_script(script=script) - - if self._config.init_script: # pragma: no cover - await self.context.add_init_script(path=self._config.init_script) - - if self._config.cookies: - await self.context.add_cookies(self._config.cookies) # pyright: ignore + if self.context: + self.context = await self._initialize_context(self._config, self.context) self._is_alive = True except Exception: @@ -340,6 +358,14 @@ class AsyncStealthySession(AsyncSession, StealthySessionMixin): else: raise RuntimeError("Session has been already started") + async def _initialize_context(self, config, ctx: AsyncBrowserContext) -> AsyncBrowserContext: + """Initialize the browser context.""" + for script in _compiled_stealth_scripts(): + await ctx.add_init_script(script=script) + + ctx = await super()._initialize_context(config, ctx) + return ctx + async def _cloudflare_solver(self, page: async_Page) -> None: # pragma: no cover """Solve the cloudflare challenge displayed on the playwright page passed @@ -443,61 +469,62 @@ class AsyncStealthySession(AsyncSession, StealthySessionMixin): ) for attempt in range(self._config.retries): - page_info = await self._get_page(params.timeout, params.extra_headers, params.disable_resources) - final_response = [None] - handle_response = self._create_response_handler(page_info, final_response) + proxy = self._config.proxy_rotator.get_proxy() if self._config.proxy_rotator else None - try: - # Navigate to URL and wait for a specified state - page_info.page.on("response", handle_response) - first_response = await page_info.page.goto(url, referer=referer) - await self._wait_for_page_stability(page_info.page, params.load_dom, params.network_idle) + async with self._page_generator( + params.timeout, params.extra_headers, params.disable_resources, proxy + ) as page_info: + final_response = [None] + page = page_info.page + page.on("response", self._create_response_handler(page_info, final_response)) - if not first_response: - raise RuntimeError(f"Failed to get response for {url}") + try: + first_response = await page.goto(url, referer=referer) + await self._wait_for_page_stability(page, params.load_dom, params.network_idle) - if params.solve_cloudflare: - await self._cloudflare_solver(page_info.page) - # Make sure the page is fully loaded after the captcha - await self._wait_for_page_stability(page_info.page, params.load_dom, params.network_idle) + if not first_response: + raise RuntimeError(f"Failed to get response for {url}") - if params.page_action: - try: - _ = await params.page_action(page_info.page) - except Exception as e: - log.error(f"Error executing page_action: {e}") + if params.solve_cloudflare: + await self._cloudflare_solver(page) + # Make sure the page is fully loaded after the captcha + await self._wait_for_page_stability(page, params.load_dom, params.network_idle) - if params.wait_selector: - try: - waiter: AsyncLocator = page_info.page.locator(params.wait_selector) - await waiter.first.wait_for(state=params.wait_selector_state) - # Wait again after waiting for the selector, helpful with protections like Cloudflare - await self._wait_for_page_stability(page_info.page, params.load_dom, params.network_idle) - except Exception as e: - log.error(f"Error waiting for selector {params.wait_selector}: {e}") + if params.page_action: + try: + _ = await params.page_action(page) + except Exception as e: # pragma: no cover + log.error(f"Error executing page_action: {e}") - await page_info.page.wait_for_timeout(params.wait) + if params.wait_selector: + try: + waiter: AsyncLocator = page.locator(params.wait_selector) + await waiter.first.wait_for(state=params.wait_selector_state) + await self._wait_for_page_stability(page, params.load_dom, params.network_idle) + except Exception as e: # pragma: no cover + log.error(f"Error waiting for selector {params.wait_selector}: {e}") - # Create response object - response = await ResponseFactory.from_async_playwright_response( - page_info.page, first_response, final_response[0], params.selector_config - ) + await page.wait_for_timeout(params.wait) - # Close the page to free up resources - await page_info.page.close() - self.page_pool.pages.remove(page_info) - return response + response = await ResponseFactory.from_async_playwright_response( + page, first_response, final_response[0], params.selector_config + ) + return response - except Exception as e: # pragma: no cover - page_info.mark_error() - await page_info.page.close() - self.page_pool.pages.remove(page_info) + except Exception as e: + page_info.mark_error() + if attempt < self._config.retries - 1: + if is_proxy_error(e): + log.warning( + f"Proxy '{proxy}' failed (attempt {attempt + 1}) | Retrying in {self._config.retry_delay}s..." + ) + else: + log.warning( + f"Attempt {attempt + 1} failed: {e}. Retrying in {self._config.retry_delay}s..." + ) + await asyncio_sleep(self._config.retry_delay) + else: + log.error(f"Failed after {self._config.retries} attempts: {e}") + raise - if attempt < self._config.retries - 1 and self._is_retriable(e): - log.warning(f"Attempt {attempt + 1} failed: {e}. Retrying in {self._config.retry_delay}s...") - await asyncio_sleep(self._config.retry_delay) - else: - raise - - # For type checking purposes only - raise AssertionError("Unreachable: retry loop must return or raise") # pragma: no cover + raise RuntimeError("Request failed") # pragma: no cover diff --git a/scrapling/engines/_browsers/_types.py b/scrapling/engines/_browsers/_types.py index 47322da..afce5d0 100644 --- a/scrapling/engines/_browsers/_types.py +++ b/scrapling/engines/_browsers/_types.py @@ -20,6 +20,7 @@ from scrapling.core._types import ( SelectorWaitStates, TYPE_CHECKING, ) +from scrapling.engines.toolbelt.proxy_rotation import ProxyRotator # Type alias for `impersonate` parameter - accepts a single browser or list of browsers ImpersonateType: TypeAlias = BrowserTypeLiteral | List[BrowserTypeLiteral] | None @@ -34,6 +35,7 @@ if TYPE_CHECKING: # pragma: no cover proxies: Optional[ProxySpec] proxy: Optional[str] proxy_auth: Optional[Tuple[str, str]] + proxy_rotator: Optional[ProxyRotator] timeout: Optional[int | float] headers: Optional[Mapping[str, Optional[str]]] retries: Optional[int] @@ -70,6 +72,7 @@ if TYPE_CHECKING: # pragma: no cover timezone_id: str | None page_action: Optional[Callable] proxy: Optional[str | Dict[str, str] | Tuple] + proxy_rotator: Optional[ProxyRotator] extra_headers: Optional[Dict[str, str]] timeout: int | float init_script: Optional[str] diff --git a/scrapling/engines/_browsers/_validators.py b/scrapling/engines/_browsers/_validators.py index 0b08f2f..2bda436 100644 --- a/scrapling/engines/_browsers/_validators.py +++ b/scrapling/engines/_browsers/_validators.py @@ -18,6 +18,7 @@ from scrapling.core._types import ( SetCookieParam, SelectorWaitStates, ) +from scrapling.engines.toolbelt.proxy_rotation import ProxyRotator from scrapling.engines.toolbelt.navigation import construct_proxy_dict from scrapling.engines._browsers._types import PlaywrightFetchParams, StealthFetchParams @@ -70,6 +71,7 @@ class PlaywrightConfig(Struct, kw_only=True, frozen=False, weakref=True): timezone_id: str | None = "" page_action: Optional[Callable] = None proxy: Optional[str | Dict[str, str] | Tuple] = None # The default value for proxy in Playwright's source is `None` + proxy_rotator: Optional[ProxyRotator] = None extra_headers: Optional[Dict[str, str]] = None timeout: Seconds = 30000 init_script: Optional[str] = None @@ -88,6 +90,11 @@ class PlaywrightConfig(Struct, kw_only=True, frozen=False, weakref=True): """Custom validation after msgspec validation""" if self.page_action and not callable(self.page_action): raise TypeError(f"page_action must be callable, got {type(self.page_action).__name__}") + if self.proxy and self.proxy_rotator: + raise ValueError( + "Cannot use 'proxy_rotator' together with 'proxy'. " + "Use either a static proxy or proxy rotation, not both." + ) if self.proxy: self.proxy = construct_proxy_dict(self.proxy) if self.cdp_url: diff --git a/scrapling/engines/constants.py b/scrapling/engines/constants.py index 9cd738e..f658416 100644 --- a/scrapling/engines/constants.py +++ b/scrapling/engines/constants.py @@ -57,7 +57,6 @@ DEFAULT_STEALTH_FLAGS = ( "--ignore-gpu-blocklist", "--enable-tcp-fast-open", "--enable-web-bluetooth", - "--disable-hang-monitor", "--disable-cloud-import", "--disable-print-preview", "--disable-dev-shm-usage", @@ -84,7 +83,6 @@ DEFAULT_STEALTH_FLAGS = ( "--prerender-from-omnibox=disabled", "--safebrowsing-disable-auto-update", "--disable-offer-upload-credit-cards", - "--disable-features=site-per-process", "--disable-background-timer-throttling", "--disable-new-content-rendering-timeout", "--run-all-compositor-stages-before-draw", diff --git a/scrapling/engines/static.py b/scrapling/engines/static.py index 46eab23..e8a41aa 100644 --- a/scrapling/engines/static.py +++ b/scrapling/engines/static.py @@ -22,9 +22,10 @@ from scrapling.core._types import ( SUPPORTED_HTTP_METHODS, ) -from ._browsers._types import RequestsSession, GetRequestParams, DataRequestParams, ImpersonateType from .toolbelt.custom import Response from .toolbelt.convertor import ResponseFactory +from .toolbelt.proxy_rotation import ProxyRotator, is_proxy_error +from ._browsers._types import RequestsSession, GetRequestParams, DataRequestParams, ImpersonateType from .toolbelt.fingerprints import generate_convincing_referer, generate_headers, __default_useragent__ _NO_SESSION: Any = object() @@ -63,6 +64,7 @@ class _ConfigurationLogic(ABC): "_default_http3", "selector_config", "_is_alive", + "_proxy_rotator", ) def __init__(self, **kwargs: Unpack[RequestsSession]): @@ -82,6 +84,13 @@ class _ConfigurationLogic(ABC): self._default_http3 = kwargs.get("http3", False) self.selector_config = kwargs.get("selector_config") or {} self._is_alive = False + self._proxy_rotator: Optional[ProxyRotator] = kwargs.get("proxy_rotator") + + if self._proxy_rotator and (self._default_proxy or self._default_proxies): + raise ValueError( + "Cannot use 'proxy_rotator' together with 'proxy' or 'proxies'. " + "Use either a static proxy or proxy rotation, not both." + ) @staticmethod def _get_param(kwargs: Dict, key: str, default: Any) -> Any: @@ -218,7 +227,7 @@ class _SyncSessionLogic(_ConfigurationLogic): selector_config = self._get_param(kwargs, "selector_config", self.selector_config) or self.selector_config max_retries = self._get_param(kwargs, "retries", self._default_retries) retry_delay = self._get_param(kwargs, "retry_delay", self._default_retry_delay) - request_args = self._merge_request_args(stealth=stealth, **kwargs) + static_proxy = kwargs.pop("proxy", None) session = self._curl_session one_off_request = False @@ -228,22 +237,38 @@ class _SyncSessionLogic(_ConfigurationLogic): session = CurlSession() one_off_request = True - if session: + if not session: + raise RuntimeError("No active session available.") # pragma: no cover + + try: for attempt in range(max_retries): + if self._proxy_rotator and static_proxy is None: + proxy = self._proxy_rotator.get_proxy() + else: + proxy = static_proxy + + request_args = self._merge_request_args(stealth=stealth, proxy=proxy, **kwargs) try: response = session.request(method, **request_args) result = ResponseFactory.from_http_request(response, selector_config) return result except CurlError as e: # pragma: no cover if attempt < max_retries - 1: - log.error(f"Attempt {attempt + 1} failed: {e}. Retrying in {retry_delay} seconds...") + # Now if the rotator is enabled, we will try again with the new proxy + # If it's not enabled, then we will try again with the same proxy + if is_proxy_error(e): + log.warning( + f"Proxy '{proxy}' failed (attempt {attempt + 1}) | Retrying in {retry_delay} seconds..." + ) + else: + log.warning(f"Attempt {attempt + 1} failed: {e}. Retrying in {retry_delay} seconds...") time_sleep(retry_delay) else: log.error(f"Failed after {max_retries} attempts: {e}") raise # Raise the exception if all retries fail - finally: - if session and one_off_request: - session.close() + finally: + if session and one_off_request: + session.close() raise RuntimeError("No active session available.") # pragma: no cover @@ -415,7 +440,7 @@ class _ASyncSessionLogic(_ConfigurationLogic): selector_config = self._get_param(kwargs, "selector_config", self.selector_config) or self.selector_config max_retries = self._get_param(kwargs, "retries", self._default_retries) retry_delay = self._get_param(kwargs, "retry_delay", self._default_retry_delay) - request_args = self._merge_request_args(stealth=stealth, **kwargs) + static_proxy = kwargs.pop("proxy", None) session = self._async_curl_session one_off_request = False @@ -427,22 +452,40 @@ class _ASyncSessionLogic(_ConfigurationLogic): session = AsyncCurlSession() one_off_request = True - if session: + if not session: + raise RuntimeError("No active session available.") # pragma: no cover + + try: + # Determine if we should use proxy rotation for attempt in range(max_retries): + if self._proxy_rotator and static_proxy is None: + proxy = self._proxy_rotator.get_proxy() + else: + proxy = static_proxy + + request_args = self._merge_request_args(stealth=stealth, proxy=proxy, **kwargs) try: response = await session.request(method, **request_args) result = ResponseFactory.from_http_request(response, selector_config) return result except CurlError as e: # pragma: no cover if attempt < max_retries - 1: - log.error(f"Attempt {attempt + 1} failed: {e}. Retrying in {retry_delay} seconds...") + # Now if the rotator is enabled, we will try again with the new proxy + # If it's not enabled, then we will try again with the same proxy + if is_proxy_error(e): + log.warning( + f"Proxy '{proxy}' failed (attempt {attempt + 1}) | Retrying in {retry_delay} seconds..." + ) + else: + log.warning(f"Attempt {attempt + 1} failed: {e}. Retrying in {retry_delay} seconds...") + await asyncio_sleep(retry_delay) else: log.error(f"Failed after {max_retries} attempts: {e}") raise # Raise the exception if all retries fail - finally: - if session and one_off_request: - await session.close() + finally: + if session and one_off_request: + await session.close() raise RuntimeError("No active session available.") # pragma: no cover @@ -604,6 +647,7 @@ class FetcherSession: "selector_config", "_client", "_is_alive", + "_proxy_rotator", ) def __init__( @@ -623,6 +667,7 @@ class FetcherSession: verify: bool = True, cert: Optional[str | Tuple[str, str]] = None, selector_config: Optional[Dict] = None, + proxy_rotator: Optional[ProxyRotator] = None, ): """ :param impersonate: Browser version to impersonate. Can be a single browser string or a list of browser strings for random selection. (Default: latest available Chrome version) @@ -641,6 +686,7 @@ class FetcherSession: :param verify: Whether to verify HTTPS certificates. Defaults to True. :param cert: Tuple of (cert, key) filenames for the client certificate. :param selector_config: Arguments passed when creating the final Selector class. + :param proxy_rotator: A ProxyRotator instance for automatic proxy rotation. """ self._default_impersonate: ImpersonateType = impersonate self._stealth = stealthy_headers @@ -659,6 +705,7 @@ class FetcherSession: self.selector_config = selector_config or {} self._is_alive = False self._client: _SyncSessionLogic | _ASyncSessionLogic | None = None + self._proxy_rotator = proxy_rotator def __enter__(self) -> _SyncSessionLogic: """Creates and returns a new synchronous Fetcher Session""" @@ -667,6 +714,7 @@ class FetcherSession: config = {k.replace("_default_", ""): getattr(self, k) for k in self.__slots__ if k.startswith("_default")} config["stealthy_headers"] = self._stealth config["selector_config"] = self.selector_config + config["proxy_rotator"] = self._proxy_rotator self._client = _SyncSessionLogic(**config) self._is_alive = True return self._client.__enter__() @@ -687,6 +735,7 @@ class FetcherSession: config = {k.replace("_default_", ""): getattr(self, k) for k in self.__slots__ if k.startswith("_default")} config["stealthy_headers"] = self._stealth config["selector_config"] = self.selector_config + config["proxy_rotator"] = self._proxy_rotator self._client = _ASyncSessionLogic(**config) self._is_alive = True return await self._client.__aenter__() diff --git a/scrapling/engines/toolbelt/__init__.py b/scrapling/engines/toolbelt/__init__.py index 8b13789..f3714c7 100644 --- a/scrapling/engines/toolbelt/__init__.py +++ b/scrapling/engines/toolbelt/__init__.py @@ -1 +1,3 @@ +from .proxy_rotation import ProxyRotator, is_proxy_error, round_robin +__all__ = ["ProxyRotator", "is_proxy_error", "round_robin"] diff --git a/scrapling/engines/toolbelt/proxy_rotation.py b/scrapling/engines/toolbelt/proxy_rotation.py new file mode 100644 index 0000000..82c270d --- /dev/null +++ b/scrapling/engines/toolbelt/proxy_rotation.py @@ -0,0 +1,104 @@ +from threading import Lock + +from scrapling.core._types import Callable, Dict, List, Tuple, ProxyType + + +RotationStrategy = Callable[[List[ProxyType], int], Tuple[ProxyType, int]] +_PROXY_ERROR_INDICATORS = { + "net::err_proxy", + "net::err_tunnel", + "connection refused", + "connection reset", + "connection timed out", + "failed to connect", + "could not resolve proxy", +} + + +def _get_proxy_key(proxy: ProxyType) -> str: + """Generate a unique key for a proxy (for dicts it's server plus username).""" + if isinstance(proxy, str): + return proxy + server = proxy.get("server", "") + username = proxy.get("username", "") + return f"{server}|{username}" + + +def is_proxy_error(error: Exception) -> bool: + """Check if an error is proxy-related. Works for both HTTP and browser errors.""" + error_msg = str(error).lower() + return any(indicator in error_msg for indicator in _PROXY_ERROR_INDICATORS) + + +def round_robin(proxies: List[ProxyType], current_index: int) -> Tuple[ProxyType, int]: + """Default round-robin rotation strategy.""" + idx = current_index % len(proxies) + return proxies[idx], (idx + 1) % len(proxies) + + +class ProxyRotator: + """ + A thread-safe proxy rotator with pluggable rotation strategies. + + Supports: + - Round-robin rotation (default) + - Custom rotation strategies via callable + - Both string URLs and Playwright-style dict proxies + """ + + __slots__ = ("_proxies", "_proxy_to_index", "_strategy", "_current_index", "_lock") + + def __init__( + self, + proxies: List[ProxyType], + strategy: RotationStrategy = round_robin, + ): + """ + Initialize the proxy rotator. + + :param proxies: List of proxy URLs or Playwright-style proxy dicts. + - String format: "http://proxy1:8080" or "http://user:pass@proxy:8080" + - Dict format: {"server": "http://proxy:8080", "username": "user", "password": "pass"} + :param strategy: Rotation strategy function. Takes (proxies, current_index) and returns (proxy, next_index). Defaults to round_robin. + """ + if not proxies: + raise ValueError("At least one proxy must be provided") + + if not callable(strategy): + raise TypeError(f"strategy must be callable, got {type(strategy).__name__}") + + self._strategy = strategy + self._lock = Lock() + + # Validate and store proxies + self._proxies: List[ProxyType] = [] + self._proxy_to_index: Dict[str, int] = {} # O(1) lookup by unique key (server + username) + for i, proxy in enumerate(proxies): + if isinstance(proxy, (str, dict)): + if isinstance(proxy, dict) and "server" not in proxy: + raise ValueError("Proxy dict must have a 'server' key") + + self._proxy_to_index[_get_proxy_key(proxy)] = i + self._proxies.append(proxy) + else: + raise TypeError(f"Invalid proxy type: {type(proxy)}. Expected str or dict.") + + self._current_index = 0 + + def get_proxy(self) -> ProxyType: + """Get the next proxy according to the rotation strategy.""" + with self._lock: + proxy, self._current_index = self._strategy(self._proxies, self._current_index) + return proxy + + @property + def proxies(self) -> List[ProxyType]: + """Get a copy of all configured proxies.""" + return list(self._proxies) + + def __len__(self) -> int: + """Return the total number of configured proxies.""" + return len(self._proxies) + + def __repr__(self) -> str: + return f"ProxyRotator(proxies={len(self._proxies)})"