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.
This commit is contained in:
@@ -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"]
|
||||
|
||||
@@ -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):
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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]
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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",
|
||||
|
||||
+62
-13
@@ -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__()
|
||||
|
||||
@@ -1 +1,3 @@
|
||||
from .proxy_rotation import ProxyRotator, is_proxy_error, round_robin
|
||||
|
||||
__all__ = ["ProxyRotator", "is_proxy_error", "round_robin"]
|
||||
|
||||
@@ -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)})"
|
||||
Reference in New Issue
Block a user