chore: перетащил обертку над браузером из scrapling

This commit is contained in:
protokey committed 2026-09-23 15:54:48 +04:00
1 parent 8f8bc62ae7
commit 257310b7d5
19 files changed
+4964 -280

No files matched your search

+1
View File
@@ -41,6 +41,7 @@ mozrunner==8.4.0
mozsystemmonitor==1.0.1 mozsystemmonitor==1.0.1
mozterm==1.0.0 mozterm==1.0.0
mozversion==2.4.0 mozversion==2.4.0
msgspec==0.21.1
multidict==6.7.1 multidict==6.7.1
numpy==2.5.3 numpy==2.5.3
orjson==3.12.0 orjson==3.12.0
View File
Whitespace-only changes.
+1 -1
View File
@@ -5,7 +5,7 @@ from patchright.async_api import Page
from antibot.base import click_and_race from antibot.base import click_and_race
from antibot.orchestrator import pass_challenges from antibot.orchestrator import pass_challenges
from captcha import load_captcha from captcha import load_captcha
from engine.schemas import ActionRequest from api.schemas import ActionRequest
_actions: dict = {} _actions: dict = {}
@@ -50,26 +50,31 @@ def load_or_create_fingerprint(user_data_dir: Path, **generator_kwargs) -> Finge
return fingerprint return fingerprint
def context_options_for(fingerprint: Fingerprint, **overrides) -> dict: def browser_kwargs_for(fingerprint: Fingerprint, browser_name: str = "chromium") -> dict:
"""Опции контекста (launch_persistent_context/new_context), согласованные с fingerprint'ом. """kwargs для `AsyncStealthySession(**kwargs)`, полностью выставляющие браузер под fingerprint —
Headers сюда намеренно не кладём — их полный набор выставляет apply_fingerprint() через UA, заголовки, init-script (подделка navigator/screen/WebGL/codecs/battery) и размер окна
set_extra_http_headers() уже после создания контекста, а он не мёржит, а заменяет, (viewport/screen/device_scale_factor).
так что выставлять их дважды бессмысленно."""
Всё уходит обычными kwargs сессии — тем же путём, что и остальная конфигурация
(proxy/cookies/locale и т.д.), как это делает сама scrapling — а не отдельным вызовом после
запуска контекста (add_init_script/set_extra_http_headers руками требуют context, который
есть только после входа в `async with AsyncStealthySession(...)`, то есть после того, как
решение о конфигурации уже "должно было" быть принято до запуска браузера).
:param fingerprint: Сгенерированный/загруженный fingerprint (см. load_or_create_fingerprint).
:param browser_name: Имя браузера для фильтрации заголовков (only_injectable_headers) —
"chromium", раз движок у нас реальный Chromium (см. docstring модуля).
:return: dict с ключами `useragent`, `extra_headers`, `init_script_content`, `additional_args`
(viewport/screen/device_scale_factor) — распаковывается в kwargs AsyncStealthySession.
"""
screen = {"width": fingerprint.screen.width, "height": fingerprint.screen.height}
return { return {
"user_agent": fingerprint.navigator.userAgent, "useragent": fingerprint.navigator.userAgent,
"viewport": { "extra_headers": only_injectable_headers(fingerprint.headers, browser_name),
"width": fingerprint.screen.width, "init_script_content": InjectFunction(fingerprint),
"height": fingerprint.screen.height, "additional_args": {
**overrides.pop("viewport", {}), "viewport": dict(screen),
}, "screen": dict(screen),
"device_scale_factor": fingerprint.screen.devicePixelRatio, "device_scale_factor": fingerprint.screen.devicePixelRatio,
**overrides, },
} }
async def apply_fingerprint(context, fingerprint: Fingerprint, browser_name: str = "chromium"):
"""Довешивает на уже созданный контекст заголовки и init-script с подделкой
navigator/screen/WebGL/codecs/battery. Вызывать сразу после launch_persistent_context/new_context,
до первого page.goto()."""
await context.set_extra_http_headers(only_injectable_headers(fingerprint.headers, browser_name))
await context.add_init_script(InjectFunction(fingerprint))
File renamed without changes.
@@ -21,7 +21,6 @@ class SolveRequest(BaseModel):
url: str url: str
proxy: str proxy: str
timeout: int | None = None timeout: int | None = None
screen: str = "1280x720"
user_agent: str | None = None user_agent: str | None = None
cookies: list[dict[str, Any]] | None = None cookies: list[dict[str, Any]] | None = None
# List[{name: str, value: str, domain: str, path: str, expires: float, httpOnly: bool, secure: bool, sameSite: Union["Lax", "None", "Strict"], partitionKey: Union[str, None]}] # List[{name: str, value: str, domain: str, path: str, expires: float, httpOnly: bool, secure: bool, sameSite: Union["Lax", "None", "Strict"], partitionKey: Union[str, None]}]
+3
View File
@@ -0,0 +1,3 @@
from engine.stealthy import AsyncStealthySession
__all__ = ["AsyncStealthySession"]
+48
View File
@@ -0,0 +1,48 @@
"""
Локальные тайп-алиасы для движка браузера — вендорено из scrapling.core._types
(урезано до того, что реально используется StealthySession/toolbelt; без RequestsSession/
GetRequestParams/DataRequestParams — они тянут curl_cffi, который в проекте не установлен
и не нужен, т.к. не используется HTTP-фетчер scrapling).
"""
from typing import (
TYPE_CHECKING,
TypeAlias,
cast,
overload,
Any,
Callable,
Dict,
Generator,
AsyncGenerator,
Generic,
List,
Set,
Literal,
Optional,
Sequence,
Tuple,
TypeVar,
Union,
Mapping,
Awaitable,
)
from typing_extensions import TypedDict
# Прокси — строка-URL или словарь в формате Playwright ({"server": ..., "username": ..., "password": ...})
ProxyType = Union[str, Dict[str, str]]
SelectorWaitStates = Literal["attached", "detached", "hidden", "visible"]
FollowRedirects = Union[bool, Literal["safe", "all", "obeycode", "firstonly"]]
# Скопировано из playwright._impl._api_structures.SetCookieParam
class SetCookieParam(TypedDict, total=False):
name: str
value: str
url: Optional[str]
domain: Optional[str]
path: Optional[str]
expires: Optional[float]
httpOnly: Optional[bool]
secure: Optional[bool]
sameSite: Optional[Literal["Lax", "None", "Strict"]]
partitionKey: Optional[str]
File diff suppressed because it is too large. Load diff
-120
View File
@@ -1,120 +0,0 @@
"""
Launch/context-флаги для Chromium, снижающие детект автоматизации на уровне ниже
JS-фингерпринта (который уже покрывает browserforge, см. engine/fingerprint.py).
Портировано из scrapling.engines.constants / _browsers._base.BaseSessionMixin.
"""
HARMFUL_ARGS = (
# Playwright сам добавляет эти флаги по умолчанию — они одни из самых дешёвых
# и надёжных сигналов автоматизации, поэтому их нужно гасить через ignore_default_args
"--enable-automation",
"--disable-popup-blocking",
"--disable-component-update",
"--disable-default-apps",
"--disable-extensions",
)
DEFAULT_ARGS = (
"--no-pings",
"--no-first-run",
"--disable-infobars",
"--disable-breakpad",
"--no-service-autorun",
"--homepage=about:blank",
"--password-store=basic",
"--disable-hang-monitor",
"--no-default-browser-check",
"--disable-session-crashed-bubble",
"--disable-search-engine-choice-screen",
)
STEALTH_ARGS = (
"--test-type",
"--mute-audio",
"--disable-sync",
"--hide-scrollbars",
"--disable-logging",
"--start-maximized", # обход headless-детекта по размеру окна
"--enable-async-dns",
"--use-mock-keychain",
"--disable-translate",
"--disable-voice-input",
"--window-position=0,0",
"--disable-wake-on-wifi",
"--ignore-gpu-blocklist",
"--enable-tcp-fast-open",
"--enable-web-bluetooth",
"--disable-cloud-import",
"--disable-print-preview",
"--disable-dev-shm-usage",
"--metrics-recording-only",
"--disable-crash-reporter",
"--disable-partial-raster",
"--disable-gesture-typing",
"--disable-checker-imaging",
"--disable-prompt-on-repost",
"--force-color-profile=srgb",
"--font-render-hinting=none",
"--aggressive-cache-discard",
"--disable-cookie-encryption",
"--disable-domain-reliability",
"--disable-threaded-animation",
"--disable-threaded-scrolling",
"--enable-simple-cache-backend",
"--disable-background-networking",
"--enable-surface-synchronization",
"--disable-image-animation-resync",
"--disable-renderer-backgrounding",
"--disable-ipc-flooding-protection",
"--prerender-from-omnibox=disabled",
"--safebrowsing-disable-auto-update",
"--disable-offer-upload-credit-cards",
"--disable-background-timer-throttling",
"--disable-new-content-rendering-timeout",
"--run-all-compositor-stages-before-draw",
"--disable-client-side-phishing-detection",
"--disable-backgrounding-occluded-windows",
"--disable-layer-tree-host-memory-pressure",
"--autoplay-policy=user-gesture-required",
"--disable-offer-store-unmasked-wallet-cards",
"--disable-blink-features=AutomationControlled",
"--disable-component-extensions-with-background-pages",
"--enable-features=NetworkService,NetworkServiceInProcess,TrustTokens,TrustTokensAlwaysAllowIssuance",
"--blink-settings=primaryHoverType=2,availableHoverTypes=2,primaryPointerType=4,availablePointerTypes=4",
"--disable-features=AudioServiceOutOfProcess,TranslateUI,BlinkGenPropertyTrees",
)
_WEBRTC_ARGS = (
"--webrtc-ip-handling-policy=disable_non_proxied_udp",
"--force-webrtc-ip-handling-policy",
)
def launch_args(locale: str | None = None, block_webrtc: bool = True) -> list[str]:
"""Флаги запуска Chromium: скорость + анти-детект + (опционально) защита от WebRTC-утечки
реального IP мимо прокси + (если передана locale) выставление языка браузера на уровне
лаунча, а не только через JS-контекст — иначе Web Workers/Intl внутри браузера остаются
на языке хоста, даже когда browserforge уже подменил navigator.language на странице."""
args = list(DEFAULT_ARGS) + list(STEALTH_ARGS)
if block_webrtc:
args += list(_WEBRTC_ARGS)
if locale:
base_lang = locale.split("-")[0].lower()
accept_lang = f"{locale},{base_lang}" if base_lang != locale.lower() else locale
args += [f"--lang={locale}", f"--accept-lang={accept_lang}"]
return args
def stealth_context_options() -> dict:
"""Context-опции, не пересекающиеся с тем, что уже выставляет browserforge
(user_agent/viewport/device_scale_factor) — их сюда специально не кладём, чтобы
fingerprint оставался единственным источником истины для этих полей."""
return {
"color_scheme": "dark", # обходит проверку prefersLightColor в creepjs
"is_mobile": False,
"has_touch": False,
"ignore_https_errors": True,
}
+97
View File
@@ -0,0 +1,97 @@
# Disable loading these resources for speed
EXTRA_RESOURCES = {
"font",
"image",
"media",
"beacon",
"object",
"imageset",
"texttrack",
"websocket",
"csp_report",
"stylesheet",
}
HARMFUL_ARGS = (
# This will be ignored to avoid detection more and possibly avoid the popup crashing bug abuse: https://issues.chromium.org/issues/340836884
"--enable-automation",
"--disable-popup-blocking",
"--disable-component-update",
"--disable-default-apps",
"--disable-extensions",
)
DEFAULT_ARGS = (
# Speed up chromium browsers by default
"--no-pings",
"--no-first-run",
"--disable-infobars",
"--disable-breakpad",
"--no-service-autorun",
"--homepage=about:blank",
"--password-store=basic",
"--disable-hang-monitor",
"--no-default-browser-check",
"--disable-session-crashed-bubble",
"--disable-search-engine-choice-screen",
)
STEALTH_ARGS = (
# Explanation: https://peter.sh/experiments/chromium-command-line-switches/
# Generally this will make the browser faster and less detectable
# "--incognito",
"--test-type",
"--mute-audio",
"--disable-sync",
"--hide-scrollbars",
"--disable-logging",
"--start-maximized", # For headless check bypass
"--enable-async-dns",
"--use-mock-keychain",
"--disable-translate",
"--disable-voice-input",
"--window-position=0,0",
"--disable-wake-on-wifi",
"--ignore-gpu-blocklist",
"--enable-tcp-fast-open",
"--enable-web-bluetooth",
"--disable-cloud-import",
"--disable-print-preview",
"--disable-dev-shm-usage",
# '--disable-popup-blocking',
"--metrics-recording-only",
"--disable-crash-reporter",
"--disable-partial-raster",
"--disable-gesture-typing",
"--disable-checker-imaging",
"--disable-prompt-on-repost",
"--force-color-profile=srgb",
"--font-render-hinting=none",
"--aggressive-cache-discard",
"--disable-cookie-encryption",
"--disable-domain-reliability",
"--disable-threaded-animation",
"--disable-threaded-scrolling",
"--enable-simple-cache-backend",
"--disable-background-networking",
"--enable-surface-synchronization",
"--disable-image-animation-resync",
"--disable-renderer-backgrounding",
"--disable-ipc-flooding-protection",
"--prerender-from-omnibox=disabled",
"--safebrowsing-disable-auto-update",
"--disable-offer-upload-credit-cards",
"--disable-background-timer-throttling",
"--disable-new-content-rendering-timeout",
"--run-all-compositor-stages-before-draw",
"--disable-client-side-phishing-detection",
"--disable-backgrounding-occluded-windows",
"--disable-layer-tree-host-memory-pressure",
"--autoplay-policy=user-gesture-required",
"--disable-offer-store-unmasked-wallet-cards",
"--disable-blink-features=AutomationControlled",
"--disable-component-extensions-with-background-pages",
"--enable-features=NetworkService,NetworkServiceInProcess,TrustTokens,TrustTokensAlwaysAllowIssuance",
"--blink-settings=primaryHoverType=2,availableHoverTypes=2,primaryPointerType=4,availablePointerTypes=4",
"--disable-features=AudioServiceOutOfProcess,TranslateUI,BlinkGenPropertyTrees",
)
+67 -32
View File
@@ -1,42 +1,35 @@
""" """
Мелкие хелперы навигации, портированные из scrapling.engines.toolbelt.navigation Блокировка ресурсов/доменов на странице и разбор строки прокси — портировано из
и scrapling.engines.toolbelt.proxy_rotation — адаптированы под patchright/async. scrapling.engines.toolbelt.navigation (async-вариант, проект целиком async).
В отличие от полного Session/PagePool из scrapling.engines._browsers._base, тут
браузер запускается и закрывается на каждый запрос (см. main.py), поэтому
переиспользовать нечего — берём только то, что относится к одной навигации:
блокировку ресурсов/доменов на странице и детект прокси-ошибок для retry.
""" """
from typing import Callable, Optional import logging
from urllib.parse import urlparse from urllib.parse import urlparse
from patchright.async_api import Route from msgspec import Struct, structs, convert, ValidationError
from playwright.async_api import Route
# ресурсы, которые можно дропать ради скорости — портировано из from engine._types import Dict, Set, Tuple, Optional, Callable
# scrapling.engines.constants.EXTRA_RESOURCES from engine.constants import EXTRA_RESOURCES
EXTRA_RESOURCES = {
"font", "image", "media", "beacon", "object",
"imageset", "texttrack", "websocket", "csp_report", "stylesheet",
}
# признаки прокси-ошибки в тексте исключения — портировано из log = logging.getLogger(__name__)
# scrapling.engines.toolbelt.proxy_rotation._PROXY_ERROR_INDICATORS
_PROXY_ERROR_INDICATORS = (
"net::err_proxy", "net::err_tunnel", "connection refused",
"connection reset", "connection timed out", "failed to connect",
"could not resolve proxy",
)
def is_proxy_error(error: Exception) -> bool: class ProxyDict(Struct):
"""Похоже ли исключение на сбой прокси (а не на обычную ошибку навигации).""" server: str
msg = str(error).lower() username: str = ""
return any(indicator in msg for indicator in _PROXY_ERROR_INDICATORS) password: str = ""
def _is_domain_blocked(hostname: str, domains: frozenset) -> bool: def _is_domain_blocked(hostname: str, domains: frozenset) -> bool:
"""Матчинг хоста и его родительских доменов за O(1) на каждый уровень — """Check if a hostname matches any blocked domain using O(1) frozenset lookups.
портировано из scrapling.engines.toolbelt.navigation._is_domain_blocked."""
Walks up the hostname's suffix chain: for "tracker.ads.doubleclick.net",
checks "tracker.ads.doubleclick.net", "ads.doubleclick.net", "doubleclick.net".
:param hostname: The hostname to check.
:param domains: A frozenset of blocked domain names.
:return: True if the hostname or any of its parent domains is in the blocked set.
"""
if hostname in domains: if hostname in domains:
return True return True
idx = hostname.find(".") idx = hostname.find(".")
@@ -48,18 +41,24 @@ def _is_domain_blocked(hostname: str, domains: frozenset) -> bool:
return False return False
def create_intercept_handler(disable_resources: bool, blocked_domains: Optional[set] = None) -> Callable: def create_async_intercept_handler(disable_resources: bool, blocked_domains: Optional[Set[str]] = None) -> Callable:
"""Обработчик route, блокирующий типы ресурсов и/или домены — портировано из """Create an async route handler that blocks both resource types and specific domains.
scrapling.engines.toolbelt.navigation.create_async_intercept_handler."""
:param disable_resources: Whether to block default resource types.
:param blocked_domains: Set of domain names to block requests to.
:return: An async route handler function.
"""
disabled_resources = EXTRA_RESOURCES if disable_resources else set() disabled_resources = EXTRA_RESOURCES if disable_resources else set()
domains = frozenset(blocked_domains) if blocked_domains else frozenset() domains = frozenset(blocked_domains) if blocked_domains else frozenset()
async def handler(route: Route) -> None: async def handler(route: Route):
if route.request.resource_type in disabled_resources: if route.request.resource_type in disabled_resources:
log.debug('Blocking background resource "%s" of type "%s"', route.request.url, route.request.resource_type)
await route.abort() await route.abort()
elif domains: elif domains:
hostname = urlparse(route.request.url).hostname or "" hostname = urlparse(route.request.url).hostname or ""
if _is_domain_blocked(hostname, domains): if _is_domain_blocked(hostname, domains):
log.debug('Blocking request to blocked domain "%s" (%s)', hostname, route.request.url)
await route.abort() await route.abort()
else: else:
await route.continue_() await route.continue_()
@@ -67,3 +66,39 @@ def create_intercept_handler(disable_resources: bool, blocked_domains: Optional[
await route.continue_() await route.continue_()
return handler return handler
def construct_proxy_dict(proxy_string: str | Dict[str, str] | Tuple) -> Dict:
"""Validate a proxy and return it in the acceptable format for Playwright
Reference: https://playwright.dev/python/docs/network#http-proxy
:param proxy_string: A string or a dictionary representation of the proxy.
:return:
"""
if isinstance(proxy_string, str):
proxy = urlparse(proxy_string)
if proxy.scheme not in ("http", "https", "socks4", "socks5") or not proxy.hostname:
raise ValueError("Invalid proxy string!")
try:
result = {
"server": f"{proxy.scheme}://{proxy.hostname}",
"username": proxy.username or "",
"password": proxy.password or "",
}
if proxy.port:
result["server"] += f":{proxy.port}"
return result
except ValueError:
# Urllib will say that one of the parameters above can't be casted to the correct type like `int` for port etc...
raise ValueError("The proxy argument's string is in invalid format!")
elif isinstance(proxy_string, dict):
try:
validated = convert(proxy_string, ProxyDict)
result_dict = structs.asdict(validated)
return result_dict
except ValidationError as e:
raise TypeError(f"Invalid proxy dictionary: {e}")
raise TypeError(f"Invalid proxy string: {proxy_string}")
+97
View File
@@ -0,0 +1,97 @@
"""
Пул страниц браузерного контекста — портировано из scrapling.engines._browsers._page.
"""
from threading import RLock
from dataclasses import dataclass
from playwright.async_api._generated import Page as AsyncPage
from engine._types import Optional, List, Literal
PageState = Literal["ready", "busy", "error"] # States that a page can be in
@dataclass
class PageInfo:
"""Information about the page and its current state"""
__slots__ = ("page", "state", "url")
page: AsyncPage
state: PageState
url: Optional[str]
def mark_busy(self, url: str = ""):
"""Mark the page as busy"""
self.state = "busy"
self.url = url
def mark_ready(self):
"""Mark the page as ready to be reused by the next request"""
self.state = "ready"
self.url = ""
def mark_error(self):
"""Mark the page as having an error"""
self.state = "error"
def __repr__(self):
return f'Page(URL="{self.url!r}", state={self.state!r})'
def __eq__(self, other_page):
"""Comparing this page to another page object."""
if other_page.__class__ is not self.__class__:
return NotImplemented
return self.page == other_page.page
class PagePool:
"""Manages a pool of browser pages/tabs with state tracking"""
__slots__ = ("max_pages", "pages", "_lock")
def __init__(self, max_pages: int = 5):
self.max_pages = max_pages
self.pages: List[PageInfo] = []
self._lock = RLock()
def add_page(self, page: AsyncPage) -> PageInfo:
"""Add a new page to the pool, marked busy for the request that created it"""
with self._lock:
if len(self.pages) >= self.max_pages:
raise RuntimeError(f"Maximum page limit ({self.max_pages}) reached")
page_info = PageInfo(page, "busy", "")
self.pages.append(page_info)
return page_info
def get_ready_page(self) -> Optional[PageInfo]:
"""Take the first ready page out of the pool's free pages, marking it busy, or return None"""
with self._lock:
for page_info in self.pages:
if page_info.state == "ready":
page_info.mark_busy()
return page_info
return None
def remove_page(self, page_info: PageInfo):
"""Forget a page, whether it's still in the pool or not"""
with self._lock:
if page_info in self.pages:
self.pages.remove(page_info)
def clear(self) -> List[PageInfo]:
"""Forget every page and return them so the caller can close them"""
with self._lock:
pages, self.pages = self.pages, []
return pages
@property
def pages_count(self) -> int:
"""Get the total number of pages"""
return len(self.pages)
@property
def busy_count(self) -> int:
"""Get the number of busy pages"""
with self._lock:
return sum(1 for p in self.pages if p.state == "busy")
+106
View File
@@ -0,0 +1,106 @@
"""
Ротация прокси и детект прокси-ошибок — портировано из scrapling.engines.toolbelt.proxy_rotation.
"""
from threading import Lock
from engine._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 cyclic_rotation(proxies: List[ProxyType], current_index: int) -> Tuple[ProxyType, int]:
"""Default cyclic rotation strategy - iterates through proxies sequentially, wrapping around at the end."""
idx = current_index % len(proxies)
return proxies[idx], (idx + 1) % len(proxies)
class ProxyRotator:
"""
A thread-safe proxy rotator with pluggable rotation strategies.
Supports:
- Cyclic 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 = cyclic_rotation,
):
"""
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 cyclic_rotation.
"""
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)})"
+391
View File
@@ -0,0 +1,391 @@
"""
Управление контекстом/пулом страниц — портировано из scrapling.engines._browsers._base.
Только async-сессия (проект целиком async, см. patchright.async_api в src/api/actions.py) и без
методов детекта Cloudflare (`_detect_cloudflare`/`_challenge_cleared` в апстриме зависят от
scrapling.parser.Selector и дублируют antibot.orchestrator.pass_challenges, который в проекте
уже решает эту задачу — см. src/api/actions.py:execute).
"""
import logging
from time import time
from re import search as re_search
from asyncio import sleep as asyncio_sleep, Lock
from contextlib import asynccontextmanager, suppress
from playwright.async_api._generated import Page as AsyncPage
from playwright.async_api import (
Frame as AsyncFrame,
Response as AsyncPlaywrightResponse,
BrowserContext as AsyncBrowserContext,
)
from playwright._impl._errors import Error as PlaywrightError
from engine.page_pool import PageInfo, PagePool
from engine.validators import validate, PlaywrightConfig, StealthConfig
from engine.navigation import construct_proxy_dict, create_async_intercept_handler
from engine._types import (
Any,
Awaitable,
Dict,
List,
Set,
Optional,
Callable,
TYPE_CHECKING,
cast,
overload,
Tuple,
ProxyType,
AsyncGenerator,
)
from engine.constants import STEALTH_ARGS, HARMFUL_ARGS, DEFAULT_ARGS
log = logging.getLogger(__name__)
class AsyncSession:
_config: "PlaywrightConfig | StealthConfig"
_context_options: Dict[str, Any]
if TYPE_CHECKING:
_build_context_with_proxy: Callable[..., Dict[str, Any]]
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: Any = None
self.context: Any = None
self.browser: Any = None
self._is_alive = False
self._lock = Lock()
async def start(self) -> None:
pass
async def close_pages(self) -> None:
"""Close every open tab in the session's pool. The next request opens a fresh tab."""
for page_info in self.page_pool.clear():
with suppress(Exception):
await cast(AsyncPage, page_info.page).close()
async def close(self):
"""Close all resources"""
if not self._is_alive: # pragma: no cover
return
await self.close_pages()
if self.context:
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
self._is_alive = False
async def __aenter__(self):
await self.start()
return self
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)
# Аналог `init_script`, но для инлайнового JS (например, InjectFunction(fingerprint) из
# browserforge — он отдаёт готовый скрипт строкой, а не путь к файлу), не апстримное
# поле scrapling.
if config.init_script_content:
await ctx.add_init_script(config.init_script_content)
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,
blocked_domains: Optional[Set[str]] = None,
context: Optional[AsyncBrowserContext] = None,
) -> PageInfo: # pragma: no cover
"""Get a ready page from the pool, or open a new one"""
ctx = context if context is not None else self.context
if TYPE_CHECKING:
assert ctx is not None, "Browser context not initialized"
async with self._lock:
page_info = self.page_pool.get_ready_page() if context is None else None
if page_info is None and context is None and self.page_pool.pages_count >= self.max_pages:
# At max capacity with the persistent context, so wait for a busy page to become ready
start_time = time()
while time() - start_time < self._max_wait_for_page:
await asyncio_sleep(0.05)
page_info = self.page_pool.get_ready_page()
if page_info is not None:
break
else:
raise TimeoutError(
f"No pages finished to clear place in the pool within the {self._max_wait_for_page}s timeout period"
)
if page_info is None:
page_info = self.page_pool.add_page(await ctx.new_page())
page = cast(AsyncPage, page_info.page)
page.set_default_navigation_timeout(timeout)
page.set_default_timeout(timeout)
await page.set_extra_http_headers(extra_headers or {})
await page.unroute_all(behavior="ignoreErrors")
if disable_resources or blocked_domains:
await page.route("**/*", create_async_intercept_handler(disable_resources, blocked_domains))
return page_info
def get_pool_stats(self) -> Dict[str, int]:
"""Get statistics about the current page pool"""
return {
"total_pages": self.page_pool.pages_count,
"busy_pages": self.page_pool.busy_count,
"max_pages": self.max_pages,
}
@staticmethod
async def _wait_for_networkidle(page: AsyncPage | AsyncFrame, timeout: Optional[int] = None):
"""Wait for the page to become idle (no network activity) even if there are never-ending requests."""
try:
await page.wait_for_load_state("networkidle", timeout=timeout)
except (PlaywrightError, Exception):
pass
async def _wait_for_page_stability(self, page: AsyncPage | AsyncFrame, load_dom: bool, network_idle: bool):
await page.wait_for_load_state(state="load")
if load_dom:
await page.wait_for_load_state(state="domcontentloaded")
if network_idle:
await self._wait_for_networkidle(page)
@staticmethod
def _create_response_handler(
page_info: PageInfo,
response_container: List,
xhr_pattern: Optional[str] = None,
xhr_container: Optional[List] = None,
) -> Callable[[AsyncPlaywrightResponse], Awaitable[None]]:
"""Create an async response handler that captures the final navigation response and optionally XHR/fetch responses.
:param page_info: The PageInfo object containing the page
:param response_container: A list to store the final response (mutable container)
:param xhr_pattern: Optional regex pattern to match XHR/fetch response URLs
:param xhr_container: Optional list to store captured XHR/fetch responses
:return: A callback function for page.on("response", ...)
"""
async def handle_response(finished_response: AsyncPlaywrightResponse) -> None:
if (
finished_response.request.resource_type == "document"
and finished_response.request.is_navigation_request()
and finished_response.request.frame == page_info.page.main_frame
):
response_container[0] = finished_response
elif (
xhr_pattern
and xhr_container is not None
and finished_response.request.resource_type in ("xhr", "fetch")
and re_search(xhr_pattern, finished_response.url)
):
xhr_container.append(finished_response)
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,
blocked_domains: Optional[Set[str]] = None,
) -> AsyncGenerator["PageInfo", None]:
"""Acquire a page - either from persistent context or fresh context with proxy."""
if proxy:
# 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)
page_info = None
try:
context = await self._initialize_context(self._config, context)
page_info = await self._get_page(
timeout, extra_headers, disable_resources, blocked_domains, context=context
)
yield page_info
finally:
if page_info is not None:
self.page_pool.remove_page(page_info)
await context.close()
else:
# Standard mode: use PagePool with persistent context
page_info = await self._get_page(timeout, extra_headers, disable_resources, blocked_domains)
try:
yield page_info
finally:
if page_info.state == "error" or page_info.page.is_closed():
with suppress(Exception):
await page_info.page.close()
self.page_pool.remove_page(page_info)
else:
page_info.mark_ready()
class BaseSessionMixin:
_config: "PlaywrightConfig | StealthConfig"
@overload
def __validate_routine__(self, params: Dict, model: type[StealthConfig]) -> StealthConfig: ...
@overload
def __validate_routine__(self, params: Dict, model: type[PlaywrightConfig]) -> PlaywrightConfig: ...
def __validate_routine__(
self, params: Dict, model: type[PlaywrightConfig] | type[StealthConfig]
) -> 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._browser_options: Dict[str, Any] = {
"args": DEFAULT_ARGS,
"ignore_default_args": HARMFUL_ARGS,
}
if "__max_pages" in params:
params["max_pages"] = params.pop("__max_pages")
config = validate(params, model=model)
self._headers_keys = (
{header.lower() for header in config.extra_headers.keys()} if config.extra_headers else set()
)
return config
def __generate_options__(self, extra_flags: Tuple | None = None) -> None:
config: PlaywrightConfig | StealthConfig = self._config
self._context_options.update(
{
"proxy": config.proxy,
"timezone_id": config.timezone_id,
"extra_http_headers": config.extra_headers,
}
)
if config.locale and config.cdp_url:
# Launch flags can't be set on remote browsers, so the detectable context option is the best effort left
self._context_options["locale"] = config.locale
if config.useragent:
self._context_options["user_agent"] = config.useragent
if not config.cdp_url:
flags = self._browser_options["args"]
if config.extra_flags or extra_flags:
flags = list(set(tuple(flags) + tuple(config.extra_flags or extra_flags or ())))
if config.dns_over_https:
doh_flag = "--dns-over-https-templates=https://cloudflare-dns.com/dns-query"
if isinstance(flags, list):
flags.append(doh_flag)
else:
flags = list(flags) + [doh_flag]
if config.locale:
# The context `locale` option patches the main thread only, so Web Workers keep the browser's real
# language and WAFs like Cloudflare flag the mismatch. Launch flags set it browser-wide instead,
# so workers, `Intl`, and the `Accept-Language` header all follow natively.
base_lang = config.locale.split("-")[0].lower()
accept_lang = f"{config.locale},{base_lang}" if base_lang != config.locale.lower() else config.locale
flags = (flags if isinstance(flags, list) else list(flags)) + [
f"--lang={config.locale}",
f"--accept-lang={accept_lang}",
]
self._browser_options.update(
{
"args": flags,
"headless": config.headless,
"channel": "chrome" if config.real_chrome else "chromium",
}
)
if config.executable_path:
self._browser_options["executable_path"] = config.executable_path
self._user_data_dir = config.user_data_dir
else:
self._browser_options = {}
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 StealthySessionMixin(BaseSessionMixin):
def __validate__(self, **params):
self._config = self.__validate_routine__(params, model=StealthConfig)
self._context_options.update(
{
"is_mobile": False,
"has_touch": False,
# I'm thinking about disabling it to rest from all Service Workers' headache, but let's keep it as it is for now
"service_workers": "allow",
"ignore_https_errors": True,
"screen": {"width": 1920, "height": 1080},
"viewport": {"width": 1920, "height": 1080},
"permissions": ["geolocation", "notifications"],
}
)
self.__generate_stealth_options()
def __generate_stealth_options(self) -> None:
config = cast(StealthConfig, self._config)
flags: Tuple[str, ...] = tuple()
if not config.cdp_url:
flags = tuple(DEFAULT_ARGS) + tuple(STEALTH_ARGS)
if config.block_webrtc:
flags += (
"--webrtc-ip-handling-policy=disable_non_proxied_udp",
"--force-webrtc-ip-handling-policy", # Ensures the policy is enforced
)
if not config.allow_webgl:
flags += (
"--disable-webgl",
"--disable-webgl-image-chromium",
"--disable-webgl2",
)
if config.hide_canvas:
flags += ("--fingerprinting-canvas-image-data-noise",)
super(StealthySessionMixin, self).__generate_options__(flags)
+80
View File
@@ -0,0 +1,80 @@
"""
TypedDict'ы для kwargs сессии и `fetch()` — портировано и урезано из
scrapling.engines._browsers._types (без RequestsSession/GetRequestParams/DataRequestParams,
которые относятся к HTTP-фетчеру scrapling и тянут curl_cffi; без `solve_cloudflare` —
см. engine/validators.py).
"""
from engine._types import (
Dict,
List,
Set,
Tuple,
Sequence,
Optional,
Callable,
TypedDict,
SetCookieParam,
SelectorWaitStates,
)
from engine.proxy_rotation import ProxyRotator
class StealthSession(TypedDict, total=False):
max_pages: int
headless: bool
disable_resources: bool
network_idle: bool
load_dom: bool
wait_selector: Optional[str]
wait_selector_state: SelectorWaitStates
cookies: Sequence[SetCookieParam] | None
google_search: bool
wait: int | float
timezone_id: str | None
page_action: Optional[Callable]
page_setup: 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]
init_script_content: Optional[str]
user_data_dir: str
selector_config: Optional[Dict]
additional_args: Optional[Dict]
locale: Optional[str]
real_chrome: bool
cdp_url: Optional[str]
useragent: Optional[str]
extra_flags: Optional[List[str]]
blocked_domains: Optional[Set[str]]
block_ads: bool
retries: int
retry_delay: int | float
capture_xhr: str | None
executable_path: Optional[str]
dns_over_https: bool
allow_webgl: bool
hide_canvas: bool
block_webrtc: bool
class PlaywrightFetchParams(TypedDict, total=False):
load_dom: bool
wait: int | float
network_idle: bool
google_search: bool
timeout: int | float
disable_resources: bool
wait_selector: Optional[str]
page_action: Optional[Callable]
page_setup: Optional[Callable]
selector_config: Optional[Dict]
extra_headers: Optional[Dict[str, str]]
wait_selector_state: SelectorWaitStates
blocked_domains: Optional[Set[str]]
proxy: Optional[str | Dict[str, str]]
class StealthFetchParams(PlaywrightFetchParams, total=False):
pass
+198
View File
@@ -0,0 +1,198 @@
"""
Stealthy-браузерная сессия (persistent context + пул страниц) — портировано и урезано из
scrapling.engines._browsers._stealth.AsyncStealthySession.
Отличия от апстрима:
- Только async-версия (проект целиком async).
- `fetch()` не строит `scrapling.Response` (это тянет `scrapling.parser.Selector` — тяжёлый
HTML-парсер, который в проекте не используется: html/cookies/user_agent достаются вручную
через сырой Playwright API). Вместо этого, если передан `page_action`, `fetch()` возвращает
РЕЗУЛЬТАТ `page_action(page)` напрямую (в апстриме он отбрасывается) — так `page_action`
становится единственной точкой, где происходит выполнение action-очереди
(`api.actions.execute`) и извлечение html/cookies/user_agent для ответа API. Если
`page_action` не передан, возвращается сам `page`.
- Без решения Cloudflare-челленджей (`_cloudflare_solver`/`solve_cloudflare`) — в проекте это
делает `antibot.orchestrator.pass_challenges`, вызываемый из `page_action`.
"""
import logging
from asyncio import sleep as asyncio_sleep
from playwright.async_api import Page, Locator
from patchright.async_api import async_playwright
from typing_extensions import Unpack
from engine.page_pool import PageInfo
from engine.proxy_rotation import is_proxy_error
from engine.session import AsyncSession, StealthySessionMixin
from engine.session_types import StealthSession, StealthFetchParams
from engine.validators import validate_fetch as _validate, StealthConfig
from engine._types import Any, List, Optional, ProxyType
log = logging.getLogger(__name__)
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.
:param headless: Run the browser in headless/hidden (default), or headful/visible mode.
:param disable_resources: Drop requests for unnecessary resources for a speed boost.
:param blocked_domains: A set of domain names to block requests to. Subdomains are also matched.
:param useragent: Pass a useragent string to be used. Otherwise the browser's own default is used.
:param cookies: Set cookies for the next request.
:param network_idle: Wait for the page until there are no network connections for at least 500 ms.
:param timeout: The timeout in milliseconds that is used in all operations and waits through the page. The default is 30,000
:param wait: The time (milliseconds) the fetcher will wait after everything finishes before closing the page and returning.
:param page_action: A function that takes the `page` object, runs after navigation, and does the automation you need
(in this project: `api.actions.execute`). Its return value is what `fetch()` returns.
:param page_setup: A function that takes the `page` object, runs before navigation. Use it to register event listeners or routes that must be set up before the page loads.
:param wait_selector: Wait for a specific CSS selector to be in a specific state.
:param init_script: An absolute path to a JavaScript file to be executed on page creation for all pages in this session.
:param locale: Specify user locale, for example, `en-GB`, `de-DE`, etc.
:param timezone_id: Changes the timezone of the browser. Defaults to the system timezone.
:param wait_selector_state: The state to wait for the selector given with `wait_selector`. The default state is `attached`.
:param real_chrome: If you have a Chrome browser installed on your device, enable this, and the Fetcher will launch an instance of your browser and use it.
:param hide_canvas: Add random noise to canvas operations to prevent fingerprinting.
:param block_webrtc: Forces WebRTC to respect proxy settings to prevent local IP address leak.
:param allow_webgl: Enabled by default. Disabling it disables WebGL and WebGL 2.0 support entirely.
:param load_dom: Enabled by default, wait for all JavaScript on page(s) to fully load and execute.
:param cdp_url: Instead of launching a new browser instance, connect to this CDP URL to control real browsers through CDP.
:param google_search: Enabled by default, Scrapling will set a Google referer header.
:param extra_headers: A dictionary of extra headers to add to the request.
:param proxy: The proxy to be used with requests, it can be a string or a dictionary with the keys 'server', 'username', and 'password' only.
:param user_data_dir: Path to a User Data Directory, which stores browser session data like cookies and local storage. The default is to create a temporary directory.
:param extra_flags: A list of additional browser flags to pass to the browser on launch.
:param additional_args: Additional arguments to be passed to Playwright's context as additional settings, and it takes higher priority than the settings above.
"""
self.__validate__(**kwargs)
super().__init__(max_pages=self._config.max_pages)
async def start(self) -> None:
"""Create a browser for this instance and context."""
if not self.playwright:
self.playwright = await async_playwright().start()
try:
if self._config.cdp_url:
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 = 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 = await self.playwright.chromium.launch_persistent_context(**persistent_options)
if self.context:
self.context = await self._initialize_context(self._config, self.context)
self._is_alive = True
except Exception:
# Clean up playwright if browser setup fails
await self.playwright.stop()
self.playwright = None
raise
else:
raise RuntimeError("Session has been already started")
async def fetch(self, url: str, **kwargs: Unpack[StealthFetchParams]) -> Any:
"""Opens up the browser and navigates to `url`.
If `page_action` is given, returns its return value (see class docstring). Otherwise returns
the Playwright `page` navigated to `url`, still open (belongs to the session's page pool/context).
:param url: The Target url.
:param google_search: Enabled by default, sets a Google referer header.
:param timeout: The timeout in milliseconds used in all operations and waits through the page.
:param wait: The time (milliseconds) to wait after everything finishes before returning.
:param page_action: A function that takes the `page` object, runs after navigation, and does the automation you need. Its return value is propagated as `fetch()`'s return value.
:param page_setup: A function that takes the `page` object, runs before navigation.
:param extra_headers: A dictionary of extra headers to add to the request.
:param disable_resources: Drop requests for unnecessary resources for a speed boost.
:param blocked_domains: A set of domain names to block requests to.
:param wait_selector: Wait for a specific CSS selector to be in a specific state.
:param wait_selector_state: The state to wait for the selector given with `wait_selector`.
:param network_idle: Wait for the page until there are no network connections for at least 500 ms.
:param load_dom: Enabled by default, wait for all JavaScript on page(s) to fully load and execute.
:param proxy: Static proxy to override rotator and session proxy. A new browser context will be created and used with it.
:return: `page_action`'s return value, or the navigated `page` if no `page_action` was given.
"""
static_proxy = kwargs.pop("proxy", None)
params = _validate(kwargs, self, StealthConfig)
if not self._is_alive: # pragma: no cover
raise RuntimeError("Context manager has been closed")
request_headers_keys = {h.lower() for h in params.extra_headers.keys()} if params.extra_headers else set()
referer = (
"https://www.google.com/" if (params.google_search and "referer" not in request_headers_keys) else None
)
for attempt in range(self._config.retries):
proxy: Optional[ProxyType] = None
if self._config.proxy_rotator and static_proxy is None:
proxy = self._config.proxy_rotator.get_proxy()
else:
proxy = static_proxy
async with self._page_generator(
params.timeout, params.extra_headers, params.disable_resources, proxy, params.blocked_domains
) as page_info:
page_info: PageInfo
page: Page = page_info.page
if params.page_setup:
try:
await params.page_setup(page)
except Exception as e: # pragma: no cover
log.error(f"Error executing page_setup: {e}")
try:
first_response = await page.goto(url, referer=referer)
await self._wait_for_page_stability(page, params.load_dom, params.network_idle)
if not first_response:
raise RuntimeError(f"Failed to get response for {url}")
result = await params.page_action(page) if params.page_action else None
if params.wait_selector:
try:
waiter: Locator = 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}")
await page.wait_for_timeout(params.wait)
return result if params.page_action else page
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
raise RuntimeError("Request failed") # pragma: no cover
+244
View File
@@ -0,0 +1,244 @@
"""
Валидация конфигурации сессии/fetch — портировано из scrapling.engines._browsers._validators.
Без поля/логики `solve_cloudflare`: детект и прохождение Cloudflare в проекте делает
antibot.orchestrator.pass_challenges (вызывается из api.actions.execute), а не сама сессия.
"""
from pathlib import Path
from typing import Annotated
from functools import lru_cache
from urllib.parse import urlparse
from dataclasses import dataclass, fields
from msgspec import Struct, Meta, convert, ValidationError
from engine._types import (
Any,
Dict,
List,
Set,
Tuple,
Optional,
Callable,
Sequence,
overload,
SetCookieParam,
SelectorWaitStates,
)
from engine.proxy_rotation import ProxyRotator
from engine.navigation import construct_proxy_dict
from engine.session_types import PlaywrightFetchParams, StealthFetchParams
# Custom validators for msgspec
@lru_cache(8)
def _is_invalid_file_path(value: str, label: str = "Init script") -> bool | str: # pragma: no cover
"""Fast file path validation"""
path = Path(value)
if not path.exists():
return f"{label} path not found: {value}"
if not path.is_file():
return f"{label} is not a file: {value}"
if not path.is_absolute():
return f"{label} is not an absolute path: {value}"
return False
@lru_cache(2)
def _is_invalid_cdp_url(cdp_url: str) -> bool | str:
"""Fast CDP URL validation"""
if not cdp_url.startswith(("ws://", "wss://", "http://", "https://")):
return "CDP URL must use 'ws://', 'wss://', 'http://', or 'https://' scheme"
netloc = urlparse(cdp_url).netloc
if not netloc: # pragma: no cover
return "Invalid hostname for the CDP URL"
return False
# Type aliases for cleaner annotations
PagesCount = Annotated[int, Meta(ge=1, le=50)]
RetriesCount = Annotated[int, Meta(ge=1, le=10)]
Seconds = Annotated[float, Meta(ge=0)]
class PlaywrightConfig(Struct, kw_only=True, frozen=False, weakref=True):
"""Configuration struct for validation"""
max_pages: PagesCount = 1
headless: bool = True
disable_resources: bool = False
network_idle: bool = False
load_dom: bool = True
wait_selector: Optional[str] = None
wait_selector_state: SelectorWaitStates = "attached"
cookies: Sequence[SetCookieParam] | None = []
google_search: bool = True
wait: Seconds = 0
timezone_id: str | None = ""
page_action: Optional[Callable] = None
page_setup: 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
init_script_content: Optional[str] = None
user_data_dir: str = ""
selector_config: Optional[Dict] = {}
additional_args: Optional[Dict] = {}
locale: str | None = None
real_chrome: bool = False
cdp_url: Optional[str] = None
useragent: Optional[str] = None
extra_flags: Optional[List[str]] = None
blocked_domains: Optional[Set[str]] = None
block_ads: bool = False
retries: RetriesCount = 3
retry_delay: Seconds = 1
capture_xhr: str | None = None
executable_path: Optional[str] = None
dns_over_https: bool = False
def __post_init__(self): # pragma: no cover
"""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.page_setup and not callable(self.page_setup):
raise TypeError(f"page_setup must be callable, got {type(self.page_setup).__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:
cdp_msg = _is_invalid_cdp_url(self.cdp_url)
if cdp_msg:
raise ValueError(cdp_msg)
if not self.cookies:
self.cookies = []
if not self.extra_flags:
self.extra_flags = []
if not self.selector_config:
self.selector_config = {}
if not self.additional_args:
self.additional_args = {}
if not self.capture_xhr:
self.capture_xhr = None
if self.init_script is not None:
validation_msg = _is_invalid_file_path(self.init_script)
if validation_msg:
raise ValueError(validation_msg)
if self.executable_path is not None:
validation_msg = _is_invalid_file_path(self.executable_path, "Browser executable")
if validation_msg:
raise ValueError(validation_msg)
if self.block_ads:
from engine.ad_domains import AD_DOMAINS
if self.blocked_domains:
self.blocked_domains = self.blocked_domains | set(AD_DOMAINS)
else:
self.blocked_domains = set(AD_DOMAINS)
class StealthConfig(PlaywrightConfig, kw_only=True, frozen=False, weakref=True):
allow_webgl: bool = True
hide_canvas: bool = False
block_webrtc: bool = False
@dataclass
class _fetch_params:
"""A dataclass of all parameters used by `fetch` calls"""
google_search: bool
timeout: Seconds
wait: Seconds
page_action: Optional[Callable]
page_setup: Optional[Callable]
extra_headers: Optional[Dict[str, str]]
disable_resources: bool
wait_selector: Optional[str]
wait_selector_state: SelectorWaitStates
network_idle: bool
load_dom: bool
blocked_domains: Optional[Set[str]]
selector_config: Dict
def validate_fetch(
method_kwargs: Dict | PlaywrightFetchParams | StealthFetchParams,
session: Any,
model: type[PlaywrightConfig] | type[StealthConfig],
) -> _fetch_params: # pragma: no cover
result: Dict[str, Any] = {}
overrides: Dict[str, Any] = {}
kwargs_dict: Dict[str, Any] = dict(method_kwargs)
# Get all field names that _fetch_params needs
fetch_param_fields = {f.name for f in fields(_fetch_params)}
for key in fetch_param_fields:
if key in kwargs_dict:
overrides[key] = kwargs_dict[key]
elif hasattr(session, "_config") and hasattr(session._config, key):
result[key] = getattr(session._config, key)
if overrides:
validated_config = validate(overrides, model)
# Extract ONLY the fields that were actually overridden (not all fields)
# This prevents validated defaults from overwriting session config values
validated_dict = {
field: getattr(validated_config, field) for field in overrides.keys() if hasattr(validated_config, field)
}
# Start with session defaults, then overwrite with validated overrides
result.update(validated_dict)
result.setdefault("blocked_domains", None)
return _fetch_params(**result)
# Cache default values for each model to reduce validation overhead
models_default_values = {}
for _model in (StealthConfig, PlaywrightConfig):
_defaults = {}
if hasattr(_model, "__struct_defaults__") and hasattr(_model, "__struct_fields__"):
for field_name, default_value in zip(_model.__struct_fields__, _model.__struct_defaults__): # type: ignore
# Skip factory defaults - these are msgspec._core.Factory instances
if type(default_value).__name__ != "Factory":
_defaults[field_name] = default_value
models_default_values[_model.__name__] = _defaults.copy()
def _filter_defaults(params: Dict, model: str) -> Dict:
"""Filter out parameters that match their default values to reduce validation overhead."""
defaults = models_default_values[model]
return {k: v for k, v in params.items() if k not in defaults or v != defaults[k]}
@overload
def validate(params: Dict, model: type[StealthConfig]) -> StealthConfig: ...
@overload
def validate(params: Dict, model: type[PlaywrightConfig]) -> PlaywrightConfig: ...
def validate(params: Dict, model: type[PlaywrightConfig] | type[StealthConfig]) -> PlaywrightConfig | StealthConfig:
try:
# Filter out params with the default values (no need to validate them) to speed up validation
filtered = _filter_defaults(params, model.__name__)
return convert(filtered, model)
except ValidationError as e:
raise TypeError(f"Invalid argument type: {e}") from e
+69 -106
View File
@@ -1,7 +1,5 @@
import asyncio
import json import json
import logging import logging
import traceback
import uuid import uuid
from pathlib import Path from pathlib import Path
from typing import Annotated from typing import Annotated
@@ -12,16 +10,13 @@ from fastapi import Cookie
from fastapi import FastAPI from fastapi import FastAPI
from fastapi import Request from fastapi import Request
from fastapi import Response from fastapi import Response
from patchright.async_api import ProxySettings from patchright.async_api import Page, ProxySettings
from patchright.async_api import async_playwright
from antibot.humanize import wait_for_page_stability from api import actions
from engine import actions from api.fingerprint import browser_kwargs_for, load_or_create_fingerprint
from engine.browser_launch import HARMFUL_ARGS, launch_args, stealth_context_options from api.geoip import resolve_geo
from engine.fingerprint import apply_fingerprint, context_options_for, load_or_create_fingerprint from api.schemas import SolveRequest
from engine.geoip import resolve_geo from engine import AsyncStealthySession
from engine.navigation import create_intercept_handler, is_proxy_error
from engine.schemas import SolveRequest
logger = logging.getLogger(__name__) logger = logging.getLogger(__name__)
logger.setLevel(logging.DEBUG) logger.setLevel(logging.DEBUG)
@@ -40,6 +35,29 @@ def parse_proxy(proxy_url: str) -> ProxySettings:
return settings return settings
def make_page_action(payload: SolveRequest):
"""Замыкание для AsyncStealthySession.fetch(page_action=...): выполняет очередь действий
(включая прохождение капчи/антибота через api.actions.execute) и достаёт из уже
навигированной страницы всё, что нужно для ответа API. Возврат этой функции — то, что
вернёт fetch() (см. engine/stealthy.py)."""
async def page_action(page: Page) -> dict:
actions_result = await actions.execute(
page, payload.actions, payload.captcha,
antibot_name=payload.captcha_type,
success_locator=payload.captcha_success_locator,
detect_locator=payload.captcha_detect_locator,
)
return {
"actions": actions_result,
"cookies": await page.context.cookies(),
"html": await page.inner_html("html"),
"user_agent": await page.evaluate("() => navigator.userAgent"),
}
return page_action
app = FastAPI() app = FastAPI()
@@ -73,20 +91,18 @@ async def solve(
response.set_cookie("session_id", session_id) response.set_cookie("session_id", session_id)
logger.info("Setup new session with session_id: %s", session_id) logger.info("Setup new session with session_id: %s", session_id)
scr_w, scr_h = payload.screen.split("x")
scr_size = int(scr_w), int(scr_h)
proxy = parse_proxy(payload.proxy) proxy = parse_proxy(payload.proxy)
# профиль (cookies/local storage/fingerprint.json) хранится на диске по # профиль (cookies/local storage/fingerprint.json) хранится на диске по
# session_id — если он уже приходил в куках, ниже подгружаем существующий # session_id — если он уже приходил в куках, ниже подгружаем существующий
# fingerprint из этой папки; если нет (новый session_id) — генерируем новый # fingerprint из этой папки; если нет (новый session_id) — генерируем новый
# и сохраняем туда же (см. engine/fingerprint.py:load_or_create_fingerprint) # и сохраняем туда же (см. api/fingerprint.py:load_or_create_fingerprint)
user_data_dir = Path(__file__).parent / f"../extra/user_data_dir/{session_id}" user_data_dir = Path(__file__).parent / f"../extra/user_data_dir/{session_id}"
if not user_data_dir.exists(): if not user_data_dir.exists():
user_data_dir.mkdir(parents=True) user_data_dir.mkdir(parents=True)
# geo — до создания fingerprint'а, чтобы производная от прокси locale # geo — до создания fingerprint'а, чтобы производная от прокси locale
# (см. engine/geoip.py) попала в сам fingerprint, а не подменяла его # (см. api/geoip.py) попала в сам fingerprint, а не подменяла его
# задним числом # задним числом
geo = await resolve_geo(proxy) geo = await resolve_geo(proxy)
locale = geo.get("locale") if geo else None locale = geo.get("locale") if geo else None
@@ -97,118 +113,65 @@ async def solve(
# TODO: если есть session_id, то user_agent нужно взять из настроек пользователя # TODO: если есть session_id, то user_agent нужно взять из настроек пользователя
if locale: if locale:
fp_kwargs["locale"] = locale fp_kwargs["locale"] = locale
scr_size = (1920, 1080)
fingerprint = load_or_create_fingerprint( fingerprint = load_or_create_fingerprint(
user_data_dir, user_data_dir,
# только нижняя граница: реальный монитор физически не может быть # только нижняя граница: реальный монитор физически не может быть
# меньше открытого на нём окна браузера (viewport ниже берём из # меньше открытого на нём окна браузера (viewport ниже берём из
# payload.screen напрямую, не отсюда). load_or_create_fingerprint # fingerprint'а, не отсюда)
# генерирует с strict=True, так что это ограничение либо честно
# соблюдается, либо сразу падает ValueError — датасет browserforge
# больше не может тихо подсунуть сюда, например, мобильный экран
screen=Screen(min_width=scr_size[0], min_height=scr_size[1]), screen=Screen(min_width=scr_size[0], min_height=scr_size[1]),
**fp_kwargs, **fp_kwargs,
) )
ctx_kwargs = { # UA/заголовки/init-script (подделка navigator/screen/WebGL/...) /viewport/screen/
# базовые анти-детект опции ниже — fingerprint накладывается поверх # device_scale_factor под fingerprint — единым куском kwargs (см.
# и выигрывает при конфликте по user_agent/viewport/device_scale_factor # api/fingerprint.py:browser_kwargs_for), тем же путём, что и остальная конфигурация
**stealth_context_options(), # сессии (proxy/cookies/locale), а не отдельным вызовом после запуска контекста.
**context_options_for( fp_browser_kwargs = browser_kwargs_for(fingerprint)
fingerprint,
viewport={"width": scr_size[0], "height": scr_size[1]},
),
}
permissions = ["notifications"] permissions = ["notifications"]
if geo: if geo and geo.get("geolocation"):
if geo.get("timezone_id"): fp_browser_kwargs["additional_args"]["geolocation"] = geo["geolocation"]
ctx_kwargs["timezone_id"] = geo["timezone_id"]
if geo.get("geolocation"):
ctx_kwargs["geolocation"] = geo["geolocation"]
permissions.append("geolocation") permissions.append("geolocation")
ctx_kwargs["permissions"] = permissions fp_browser_kwargs["additional_args"]["permissions"] = permissions
nav_timeout_ms = payload.timeout * 1000 if payload.timeout else 60_000 nav_timeout_ms = payload.timeout * 1000 if payload.timeout else 60_000
default_timeout_ms = payload.timeout * 1000 if payload.timeout else 30_000
blocked_domains = set(payload.blocked_domains) if payload.blocked_domains else None blocked_domains = set(payload.blocked_domains) if payload.blocked_domains else None
referer = "https://www.google.com/" if payload.google_search else None
retries = max(1, payload.retries)
async with async_playwright() as p: kwargs = dict(
context = await p.chromium.launch_persistent_context( headless=False,
user_data_dir=str(user_data_dir),
proxy=proxy, proxy=proxy,
headless=True, user_data_dir=str(user_data_dir),
channel="chrome", real_chrome=True,
args=launch_args(locale=locale), locale=locale,
ignore_default_args=HARMFUL_ARGS, timezone_id=geo.get("timezone_id") if geo else None,
**ctx_kwargs, cookies=payload.cookies,
disable_resources=payload.disable_resources,
blocked_domains=blocked_domains,
timeout=nav_timeout_ms,
retries=max(1, payload.retries),
retry_delay=payload.retry_delay,
wait=payload.wait,
google_search=payload.google_search
) )
try:
await apply_fingerprint(context, fingerprint)
page = context.pages[0] if context.pages else await context.new_page()
if payload.cookies:
await context.add_cookies(payload.cookies)
page.set_default_timeout(default_timeout_ms)
page.set_default_navigation_timeout(nav_timeout_ms)
if payload.disable_resources or blocked_domains:
await page.route("**/*", create_intercept_handler(payload.disable_resources, blocked_domains))
status = error = user_agent = None async with AsyncStealthySession(**kwargs) as engine:
cookies = html = None result = await engine.fetch(payload.url, page_action=make_page_action(payload))
actions_result = None
for attempt in range(retries):
try:
await page.goto(url=payload.url, referer=referer)
await wait_for_page_stability(page, load_dom=True, network_idle=False)
actions_result = await actions.execute(
page, payload.actions, payload.captcha,
antibot_name=payload.captcha_type,
success_locator=payload.captcha_success_locator,
detect_locator=payload.captcha_detect_locator,
)
if payload.wait:
await page.wait_for_timeout(payload.wait)
status = not "error" in actions_result[-1]
error = ["Something error", None][status]
cookies = await context.cookies()
html = await page.inner_html('html')
user_agent = await page.evaluate("() => navigator.userAgent")
# TODO: добавить как опцию. передачу скрина в base64.
# await page.screenshot(
# path=(
# f"/home/sokol/PycharmProjects/WebRoboApi/extra/screenshots/"
# f"{datetime.timestamp(datetime.now())}.png"
# )
# )
break
except Exception as e:
if attempt < retries - 1:
kind = "прокси" if is_proxy_error(e) else "запрос"
logger.warning(
"Попытка %s/%s не удалась (%s: %s), retry через %.1fs",
attempt + 1, retries, kind, e, payload.retry_delay,
)
await asyncio.sleep(payload.retry_delay)
else:
raise
finally:
await context.close()
finally: finally:
BROWSER_SLOTS += 1 BROWSER_SLOTS += 1
result = { status = not "error" in result["actions"][-1]
error = ["Something error", None][status]
result_payload = {
"status": ["err", "ok"][status], "status": ["err", "ok"][status],
"user_agent": user_agent, "user_agent": result["user_agent"],
"cookies": cookies, "cookies": result["cookies"],
"actions": actions_result, "actions": result["actions"],
} }
if payload.html: if payload.html:
result["html"] = html result_payload["html"] = result["html"]
if error: if error:
result["error"] = error result_payload["error"] = error
logger.info("Outgoing payload:\n%s", json.dumps(result, indent=2, sort_keys=True, ensure_ascii=True)) logger.info("Outgoing payload:\n%s", json.dumps(result_payload, indent=2, sort_keys=True, ensure_ascii=True))
return result return result_payload