diff --git a/docs/topics/throttling.rst b/docs/topics/throttling.rst index 852b9f8e9..436d1e336 100644 --- a/docs/topics/throttling.rst +++ b/docs/topics/throttling.rst @@ -45,6 +45,10 @@ The main throttling :ref:`settings ` are: Even if you have multiple slots, requests to the same domain cannot be sent more frequently than this delay. + To target a specific number of requests per minute (RPM) *per domain*, set + this to ``60 / RPM``. For example, ``DOWNLOAD_DELAY = 1.0`` for 60 RPM, or + ``DOWNLOAD_DELAY = 2.0`` for 30 RPM. + - .. setting:: DOWNLOAD_DELAY_PER_SLOT :setting:`DOWNLOAD_DELAY_PER_SLOT` (default: ``None``) @@ -945,14 +949,6 @@ Additional settings If ``True``, ``0.5`` (i.e. ±50%) is used as the randomization factor. If ``False``, no randomization is applied. -- .. setting:: TARGET_RPM - - :setting:`TARGET_RPM` (default: ``None``) - - Target number of requests per minute *per domain*. It has no effect on its - own; it is read by :class:`~scrapy.addons.throttling.TargetRPMAddon`, which - must be enabled explicitly. - - .. setting:: THROTTLING_DEBUG :setting:`THROTTLING_DEBUG` (default: ``False``) diff --git a/scrapy/settings/default_settings.py b/scrapy/settings/default_settings.py index 709d55038..f19894658 100644 --- a/scrapy/settings/default_settings.py +++ b/scrapy/settings/default_settings.py @@ -584,8 +584,6 @@ STATS_DUMP = True STATSMAILER_RCPTS = [] -TARGET_RPM = None - TELNETCONSOLE_ENABLED = 1 TELNETCONSOLE_HOST = "127.0.0.1" TELNETCONSOLE_PORT = [6023, 6073] diff --git a/scrapy/throttling.py b/scrapy/throttling.py index 1f5168cd5..9b70ec23f 100644 --- a/scrapy/throttling.py +++ b/scrapy/throttling.py @@ -1,29 +1,38 @@ from __future__ import annotations +import contextlib import datetime as dt -from collections.abc import Awaitable, Iterable -from datetime import UTC +import logging +import random +import time +from collections.abc import Awaitable, Callable, Iterable from email.utils import parsedate_to_datetime from functools import wraps -from typing import TYPE_CHECKING, Any, Callable, Protocol, TypedDict, Union +from typing import TYPE_CHECKING, Any, Protocol, TypedDict, TypeVar, cast from weakref import WeakKeyDictionary +from twisted.internet.defer import Deferred from typing_extensions import NotRequired, Self -from scrapy.http import Request, Response +from scrapy import signals +from scrapy.utils.asyncio import sleep, wait_for_first from scrapy.utils.httpobj import urlparse_cached -from scrapy.utils.misc import load_object +from scrapy.utils.misc import build_from_crawler, load_object if TYPE_CHECKING: from scrapy.crawler import Crawler + from scrapy.http import Request, Response + + +logger = logging.getLogger(__name__) def _parse_retry_after(response: Response) -> float | None: - value = response.headers.get("Retry-After") - if not value: + raw = response.headers.get("Retry-After") + if not raw: return None try: - value = value.decode("utf-8").strip() + value = raw.decode("utf-8").strip() except UnicodeDecodeError: return None if value.isdigit(): @@ -33,18 +42,18 @@ def _parse_retry_after(response: Response) -> float | None: except (TypeError, ValueError, OverflowError): return None if date.tzinfo is None: - date = date.replace(tzinfo=UTC) - now = dt.datetime.now(UTC) + date = date.replace(tzinfo=dt.timezone.utc) + now = dt.datetime.now(dt.timezone.utc) seconds_to_wait = (date - now).total_seconds() return max(0, int(seconds_to_wait)) or None def _parse_ratelimit_reset(response: Response) -> float | None: - value = response.headers.get("RateLimit-Reset") - if not value: + raw = response.headers.get("RateLimit-Reset") + if not raw: return None try: - value = value.decode("utf-8").strip() + value = raw.decode("utf-8").strip() except UnicodeDecodeError: return None try: @@ -59,9 +68,50 @@ class BackoffScopeData(TypedDict): remaining: NotRequired[float] +class BackoffConfig(TypedDict, total=False): + """Per-scope override of the backoff settings. + + Used as the value of the ``"backoff"`` key of :class:`ThrottlingScopeConfig` + entries. Any key left out falls back to the corresponding global + ``BACKOFF_*`` setting. + """ + + http_codes: list[int] + exceptions: list[str] + delay_factor: float + max_delay: float + min_delay: float + jitter: float | list[float] + + +class ThrottlingScopeConfig(TypedDict, total=False): + """Accepted keys of :setting:`THROTTLING_SCOPES` entries. + + Every key is optional; missing keys fall back to the matching global + setting (e.g. ``delay`` falls back to :setting:`DOWNLOAD_DELAY`). + """ + + concurrency: int + + min_concurrency: int + """Floor. Never drop below this during backoff/rampup.""" + + delay: float + jitter: float | list[float] + quota: float + window: float + rampup: bool + + manager: str | type + """Import path or class of a custom :setting:`THROTTLING_SCOPE_MANAGER` for + this scope.""" + + backoff: BackoffConfig + + ScopeID = str -BackoffData = Union[None, ScopeID, Iterable[ScopeID], dict[ScopeID, BackoffScopeData]] -RequestScopes = Union[None, ScopeID, Iterable[ScopeID], dict[ScopeID, float | None]] +BackoffData = None | ScopeID | Iterable[ScopeID] | dict[ScopeID, BackoffScopeData] +RequestScopes = None | ScopeID | Iterable[ScopeID] | dict[ScopeID, float | None] def iter_scopes(scopes: RequestScopes) -> Iterable[ScopeID]: @@ -74,6 +124,61 @@ def iter_scopes(scopes: RequestScopes) -> Iterable[ScopeID]: return iter(scopes) +def iter_scope_values(scopes: RequestScopes) -> Iterable[tuple[ScopeID, float | None]]: + """Iterate over *scopes* as ``(scope_id, value)`` pairs. + + For dict scopes the value is the expected :ref:`throttling quota + ` consumption; for every other form the value is + ``None``. + """ + if scopes is None: + return + if isinstance(scopes, str): + yield scopes, None + return + if isinstance(scopes, dict): + yield from scopes.items() + return + for scope in scopes: + yield scope, None + + +def _to_scope_dict(collection: Any, default: Callable[[], Any]) -> dict[ScopeID, Any]: + """Normalize *collection* (``None``, str, iterable or dict) into a dict + mapping scope names to values produced by *default*.""" + if isinstance(collection, dict): + return collection + if collection is None: + return {} + if isinstance(collection, str): + return {collection: default()} + if isinstance(collection, Iterable): + return {scope: default() for scope in collection} + raise TypeError( + f"Invalid type ({type(collection)}) of scopes value " + f"{collection!r}. Expected None, str, Iterable or dict." + ) + + +def _add_bare_scope(collection: Any, scope: ScopeID, empty: Any) -> Any: + """Add *scope* to *collection* without any associated value, keeping the + most compact representation possible.""" + if collection is None: + return scope + if isinstance(collection, str): + return collection if collection == scope else {collection, scope} + if isinstance(collection, dict): + if scope not in collection: + collection[scope] = empty + return collection + if isinstance(collection, Iterable): + return set(collection) | {scope} if scope not in collection else collection + raise TypeError( + f"Invalid type ({type(collection)}) of scopes value " + f"{collection!r}. Expected None, str, Iterable or dict." + ) + + def add_scope( scopes: RequestScopes, scope: ScopeID, @@ -86,38 +191,12 @@ def add_scope( :meth:`~ThrottlingManagerProtocol.get_scopes`, e.g. in :class:`ThrottlingManager` subclasses. """ - if value is not None: - if not isinstance(scopes, dict): - if scopes is None: - scopes = {} - elif isinstance(scopes, str): - scopes = {scopes: None} - elif isinstance(scopes, Iterable): - scopes = {s: None for s in scopes} - else: - raise TypeError( - f"Invalid type ({type(scopes)}) of scopes value " - f"{scopes!r}. Expected None, str, Iterable or dict." - ) - if scope in scopes and not isinstance(scopes[scope], dict): - raise TypeError(f"Scope {scope!r} has a non-dict value in {scopes!r}") - scopes[scope] = value - elif scopes is None: - scopes = scope - elif isinstance(scopes, str): - if scopes != scope: - scopes = {scopes, scope} - elif isinstance(scopes, dict): - if scope not in scopes: - scopes[scope] = None - elif isinstance(scopes, Iterable): - if scope not in scopes: - scopes = set(scopes) | {scope} - else: - raise TypeError( - f"Invalid type ({type(scopes)}) of scopes value " - f"{scopes!r}. Expected None, str, Iterable or dict." - ) + if value is None: + return cast("RequestScopes", _add_bare_scope(scopes, scope, None)) + scopes = _to_scope_dict(scopes, lambda: None) + if scope in scopes and not isinstance(scopes[scope], dict): + raise TypeError(f"Scope {scope!r} has a non-dict value in {scopes!r}") + scopes[scope] = value return scopes @@ -129,7 +208,7 @@ def update_scope_backoff( delay: float | None = None, consumed: float | None = None, ) -> BackoffData: - """Add *scope* to *backoff* or update its existing entry the given + """Add *scope* to *backoff* or update its existing entry with the given parameters. This is a utility function to help extending the output of @@ -138,45 +217,16 @@ def update_scope_backoff( :meth:`~ThrottlingManagerProtocol.get_exception_backoff`, e.g. in :class:`ThrottlingManager` subclasses. """ - has_params = delay is not None or consumed is not None - if has_params: - if not isinstance(backoff, dict): - if backoff is None: - backoff = {} - elif isinstance(backoff, str): - backoff = {backoff: {}} - elif isinstance(backoff, Iterable): - backoff = {s: {} for s in backoff} - else: - raise TypeError( - f"Invalid type ({type(backoff)}) of scopes value " - f"{backoff!r}. Expected None, str, Iterable or dict." - ) - if scope in backoff: - if not isinstance(backoff[scope], dict): - raise TypeError(f"Scope {scope!r} has a non-dict value in {backoff!r}") - else: - backoff[scope] = {} - if delay is not None: - backoff[scope]["delay"] = delay - if consumed is not None: - backoff[scope]["consumed"] = consumed - elif backoff is None: - backoff = scope - elif isinstance(backoff, str): - if backoff != scope: - backoff = {backoff, scope} - elif isinstance(backoff, dict): - if scope not in backoff: - backoff[scope] = {} - elif isinstance(backoff, Iterable): - if scope not in backoff: - backoff = set(backoff) | {scope} - else: - raise TypeError( - f"Invalid type ({type(backoff)}) of scopes value " - f"{backoff!r}. Expected None, str, Iterable or dict." - ) + if delay is None and consumed is None: + return cast("BackoffData", _add_bare_scope(backoff, scope, {})) + backoff = _to_scope_dict(backoff, dict) + entry = backoff.setdefault(scope, {}) + if not isinstance(entry, dict): + raise TypeError(f"Scope {scope!r} has a non-dict value in {backoff!r}") + if delay is not None: + entry["delay"] = delay + if consumed is not None: + entry["consumed"] = consumed return backoff @@ -251,13 +301,35 @@ class ThrottlingManagerProtocol(Protocol): It supports the same return values as :meth:`get_initial_backoff`. """ + async def acquire(self, request: Request) -> None: + """Block until *request* is allowed to be sent by all of its scopes. -GetScopesMethod = Callable[ - [ThrottlingManagerProtocol, Request], Awaitable[RequestScopes] -] + This is the throttling gate that the engine awaits before releasing a + request to the downloader. + """ + + def release(self, request: Request) -> None: + """Release the concurrency slots that :meth:`acquire` reserved for + *request*. + + The engine calls this once *request* has finished downloading (whether + it succeeded, failed or returned a new request), so that scopes that + enforce a concurrency limit can let other requests through. + """ + + async def process_response(self, response: Response) -> None: + """Update the throttling state based on *response*.""" + + async def process_exception(self, request: Request, exception: Exception) -> None: + """Update the throttling state based on a download *exception*.""" -def scope_cache(f: GetScopesMethod) -> GetScopesMethod: +_GetScopesMethod = TypeVar( + "_GetScopesMethod", bound=Callable[..., Awaitable[RequestScopes]] +) + + +def scope_cache(f: _GetScopesMethod) -> _GetScopesMethod: """Decorator to cache the result of :meth:`~ThrottlingManagerProtocol.get_scopes` calls. @@ -280,17 +352,17 @@ def scope_cache(f: GetScopesMethod) -> GetScopesMethod: async def get_scopes(self, request): return urlparse_cached(request).netloc """ - cache = WeakKeyDictionary() + cache: WeakKeyDictionary[Request, RequestScopes] = WeakKeyDictionary() @wraps(f) - async def wrapper(self, request: Request): + async def wrapper(self: Any, request: Request) -> RequestScopes: if request in cache: return cache[request] scopes = await f(self, request) cache[request] = scopes return scopes - return wrapper + return wrapper # type: ignore[return-value] class ThrottlingManager: @@ -306,27 +378,57 @@ class ThrottlingManager: def __init__(self, crawler: Crawler) -> None: self.crawler = crawler - self.throttler = crawler.throttler - self.backoff_http_codes = set(crawler.settings.getlist("BACKOFF_HTTP_CODES")) - self.backoff_exceptions = tuple( + self._backoff_http_codes = { + int(code) for code in crawler.settings.getlist("BACKOFF_HTTP_CODES") + } + self._backoff_exceptions = tuple( load_object(cls) for cls in crawler.settings.getlist("BACKOFF_EXCEPTIONS") ) + self._debug = crawler.settings.getbool("THROTTLING_DEBUG") + self._max_idle = crawler.settings.getfloat("THROTTLING_SCOPE_MAX_IDLE") + self._robotstxt_obey = crawler.settings.getbool( + "ROBOTSTXT_OBEY" + ) and crawler.settings.getbool("THROTTLING_ROBOTSTXT_OBEY") + self._robotstxt_max_delay = crawler.settings.getfloat( + "THROTTLING_ROBOTSTXT_MAX_DELAY" + ) + self._default_useragent: str = crawler.settings["USER_AGENT"] + self._robotstxt_useragent: str | None = crawler.settings["ROBOTSTXT_USER_AGENT"] + if self._robotstxt_obey: + crawler.signals.connect( + self._on_robots_parsed, signal=signals.robots_parsed + ) + self._default_scope_manager_cls = load_object( + crawler.settings["THROTTLING_SCOPE_MANAGER"] + ) + self._scopes_config: dict[str, dict[str, Any]] = crawler.settings.getdict( + "THROTTLING_SCOPES" + ) + self._scope_managers: dict[ScopeID, ThrottlingScopeManagerProtocol] = {} + self._last_eviction: float | None = None + # Concurrency slots reserved by acquire(), to be released once the + # request finishes downloading. + self._reserved: WeakKeyDictionary[ + Request, list[tuple[ThrottlingScopeManagerProtocol, float | None]] + ] = WeakKeyDictionary() @scope_cache - async def get_scopes( - self: ThrottlingManagerProtocol, request: Request - ) -> RequestScopes: + async def get_scopes(self, request: Request) -> RequestScopes: + scopes = request.meta.get("throttling_scopes") + if scopes is not None: + return cast("RequestScopes", scopes) return urlparse_cached(request).netloc async def get_initial_backoff(self) -> BackoffData: return None async def get_response_backoff(self, response: Response) -> BackoffData: - if response.status not in self.backoff_http_codes: - return None assert response.request is not None - assert self.throttler is not None - scopes = await self.throttler.get_scopes(response.request) + if response.request.meta.get("throttling_dont_track"): + return None + if response.status not in self._backoff_http_codes: + return None + scopes = await self.get_scopes(response.request) if delay := self.get_response_delay(response): scopes = {scope: {"delay": delay} for scope in iter_scopes(scopes)} return scopes @@ -347,11 +449,207 @@ class ThrottlingManager: async def get_exception_backoff( self, request: Request, exception: Exception ) -> BackoffData: - if isinstance(exception, self.backoff_exceptions): - assert self.throttler is not None - return await self.throttler.get_scopes(request) + if request.meta.get("throttling_dont_track"): + return None + if isinstance(exception, self._backoff_exceptions): + return await self.get_scopes(request) return None + # -- Scope-state coordination (called from the request lifecycle) -------- + + def _get_scope_manager(self, scope_id: ScopeID) -> ThrottlingScopeManagerProtocol: + manager = self._scope_managers.get(scope_id) + if manager is None: + config: dict[str, Any] = dict(self._scopes_config.get(scope_id, {})) + config.setdefault("id", scope_id) + manager_cls = ( + load_object(config["manager"]) + if "manager" in config + else self._default_scope_manager_cls + ) + manager = build_from_crawler(manager_cls, self.crawler, config) + self._scope_managers[scope_id] = manager + return manager + + async def acquire(self, request: Request) -> None: + now = time.monotonic() + self._maybe_evict(now) + await self._apply_request_delay(request) + scope_values = list(iter_scope_values(await self.get_scopes(request))) + if not scope_values: + return + managers = [ + (self._get_scope_manager(scope_id), value) + for scope_id, value in scope_values + ] + while True: + wait = max( + [0.0, *(manager.can_send(amount=value) for manager, value in managers)] + ) + if wait > 0: + if self._debug: + logger.debug( + f"Throttling {request} for {wait:.2f}s " + f"(scopes: {[scope_id for scope_id, _ in scope_values]})" + ) + await sleep(wait) + continue + # All time-based gates (delay, backoff, quota) are open; the only + # remaining reason to wait is a full concurrency slot. + blocked = [ + manager for manager, _ in managers if manager.concurrency_blocked() + ] + if not blocked: + for manager, value in managers: + manager.record_sent(amount=value) + self._reserved[request] = managers + return + if self._debug: + logger.debug( + f"Throttling {request} until a concurrency slot frees up " + f"(scopes: {[scope_id for scope_id, _ in scope_values]})" + ) + await self._wait_for_slot(blocked) + + def release(self, request: Request) -> None: + managers = self._reserved.pop(request, None) + if not managers: + return + for manager, _ in managers: + manager.record_done() + + async def _wait_for_slot(self, managers: list[Any]) -> None: + """Block until any of *managers* frees a concurrency slot. + + Each manager hands out an event Deferred that fires when a slot is freed + (via :meth:`ThrottlingScopeManager.record_done`) or the limit is raised + (via :meth:`ThrottlingScopeManager.set_concurrency`). A long safety timer + bounds the wait in case no slot is ever freed (it always should be, via + :meth:`release`). + """ + pairs = [(manager, manager.slot_event()) for manager in managers] + events = [event for _, event in pairs] + _, pending = await wait_for_first(events, timeout=_SLOT_WAIT_TIMEOUT) + for manager, event in pairs: + if event in pending: + manager.discard_slot_event(event) + + async def _apply_request_delay(self, request: Request) -> None: + """Honor the :reqmeta:`throttling_delay` meta key by holding *request* + for the requested number of seconds the first time it is processed.""" + delay = request.meta.get("throttling_delay") + if not delay or request.meta.get("_throttling_delayed"): + return + request.meta["_throttling_delayed"] = True + if self._debug: + logger.debug(f"Holding {request} for {delay:.2f}s (throttling_delay)") + await sleep(float(delay)) + + async def process_response(self, response: Response) -> None: + data = await self.get_response_backoff(response) + self._apply_backoff(data) + + async def process_exception(self, request: Request, exception: Exception) -> None: + data = await self.get_exception_backoff(request, exception) + self._apply_backoff(data) + + def _apply_backoff(self, data: BackoffData) -> None: + if data is None: + return + if isinstance(data, dict): + items: Iterable[tuple[ScopeID, Any]] = data.items() + else: + items = ((scope_id, None) for scope_id in iter_scopes(data)) + for scope_id, entry in items: + manager = self._get_scope_manager(scope_id) + delay = consumed = remaining = None + if isinstance(entry, dict): + delay = entry.get("delay") + consumed = entry.get("consumed") + remaining = entry.get("remaining") + # A dict entry that only reports quota usage reconciles the quota + # without counting as a backoff trigger. + quota_only = ( + isinstance(entry, dict) + and delay is None + and (consumed is not None or remaining is not None) + ) + if consumed is not None or remaining is not None: + manager.reconcile_quota(consumed=consumed, remaining=remaining) + if quota_only: + continue + if self._debug: + logger.debug(f"Backoff for scope {scope_id} (delay: {delay})") + manager.record_backoff(delay=delay) + + def _on_robots_parsed(self, robotparser: Any, request: Request) -> None: + """Honor a robots.txt ``Crawl-delay`` on the :signal:`robots_parsed` + signal. + + It reads the ``Crawl-delay`` directive for the configured user agent from + the parsed robots.txt and, if present, applies it to the scope of the + host that *request* targets via :meth:`apply_robots_crawl_delay`. + """ + if not self._robotstxt_obey: + return + useragent: str | bytes = self._robotstxt_useragent or self._default_useragent + try: + delay = robotparser.crawl_delay(useragent) + except Exception: # pragma: no cover - backend-specific failures + return + if delay: + self.apply_robots_crawl_delay(urlparse_cached(request).netloc, delay) + + def apply_robots_crawl_delay(self, scope_id: ScopeID, delay: float) -> None: + """Honor a robots.txt ``Crawl-delay`` directive of *delay* seconds for + *scope_id* by setting its delay (capped at + :setting:`THROTTLING_ROBOTSTXT_MAX_DELAY`) and its concurrency to ``1``. + + Called from the :signal:`robots_parsed` signal handler when + :setting:`THROTTLING_ROBOTSTXT_OBEY` is enabled. An explicit + :setting:`THROTTLING_SCOPES` configuration for the scope is respected, but + a warning is logged about the discrepancy unless its ``ignore_robots_txt`` + key is ``True``. + """ + if not self._robotstxt_obey: + return + capped = min(delay, self._robotstxt_max_delay) + config = self._scopes_config.get(scope_id, {}) + if config.get("ignore_robots_txt"): + return + conflicts = [] + if config.get("delay") is not None and float(config["delay"]) < capped: + conflicts.append(f"delay={config['delay']!r} < Crawl-delay {capped}") + if config.get("concurrency") is not None and int(config["concurrency"]) > 1: + conflicts.append(f"concurrency={config['concurrency']!r} > 1") + if conflicts: + logger.warning( + f"Throttling scope {scope_id!r} is configured with {' and '.join(conflicts)}, " + f"which is more aggressive than its robots.txt Crawl-delay of " + f"{capped}s. The configured values take precedence; set " + "'ignore_robots_txt': True in its THROTTLING_SCOPES entry to " + "silence this warning." + ) + return + if self._debug: + logger.debug(f"robots.txt Crawl-delay for scope {scope_id}: {capped}s") + manager = self._get_scope_manager(scope_id) + manager.set_base_delay(capped) + manager.set_concurrency(1) + + def _maybe_evict(self, now: float) -> None: + if self._max_idle <= 0: + return + if ( + self._last_eviction is not None + and now - self._last_eviction < self._max_idle / 2 + ): + return + self._last_eviction = now + for scope_id in list(self._scope_managers): + if self._scope_managers[scope_id].is_idle(now, self._max_idle): + del self._scope_managers[scope_id] + class ThrottlingScopeManagerProtocol(Protocol): """A protocol for :setting:`THROTTLING_SCOPE_MANAGER` :ref:`components @@ -364,7 +662,7 @@ class ThrottlingScopeManagerProtocol(Protocol): { "id": "example.com", - "concurrency": 1.0, + "concurrency": 1, "delay": 1.0, "jitter": 0.5, "quota": 1000.0, @@ -376,7 +674,6 @@ class ThrottlingScopeManagerProtocol(Protocol): "max_delay": 180.0, "min_delay": 5.0, "jitter": [0.01, 0.33], - "concurrency_factor": 0.8, }, "rampup": { "backoff_target": 1, @@ -384,6 +681,7 @@ class ThrottlingScopeManagerProtocol(Protocol): "min_delay": 0.05, }, } + """ @classmethod @@ -393,6 +691,378 @@ class ThrottlingScopeManagerProtocol(Protocol): def __init__(self, crawler: Crawler, config: dict[str, Any]) -> None: pass + def can_send(self, now: float | None = None, amount: float | None = None) -> float: + """Return the number of seconds to wait before a request for this scope + may be sent, or ``0`` if it may be sent right away. + + *amount* is the expected :ref:`throttling quota ` + consumption of the request, if any. + """ + + def record_sent( + self, now: float | None = None, amount: float | None = None + ) -> None: + """Record that a request for this scope has just been sent, consuming + *amount* of its :ref:`throttling quota ` if given.""" + + def record_done(self, now: float | None = None) -> None: + """Record that a previously :meth:`record_sent` request has finished + downloading, freeing its concurrency slot.""" + + def record_backoff( + self, delay: float | None = None, now: float | None = None + ) -> None: + """Apply a backoff to this scope. + + *delay*, when given, is a hard minimum delay in seconds (e.g. from a + ``Retry-After`` header). When omitted, an exponential backoff step is + applied instead. + """ + + def reconcile_quota( + self, + consumed: float | None = None, + remaining: float | None = None, + now: float | None = None, + ) -> None: + """Reconcile the :ref:`throttling quota ` of this + scope with the actual *consumed* amount (or the *remaining* amount) + reported for a request, correcting the estimate used by + :meth:`record_sent`.""" + + def set_base_delay(self, delay: float) -> None: + """Raise the base (non-backoff) delay of this scope to *delay* seconds. + + It never lowers the configured base delay; it is used to honor external + hints such as a robots.txt ``Crawl-delay`` directive. + """ + + def set_concurrency(self, concurrency: int) -> None: + """Set the maximum number of concurrent requests allowed for this + scope.""" + + def concurrency_blocked(self) -> bool: + """Return whether this scope is at its concurrency limit. + + :class:`ThrottlingManager` calls this (after all time-based gates in + :meth:`can_send` are open) to decide whether to wait for a freed slot. + Return ``False`` when no concurrency limit is enforced. + """ + + def slot_event(self) -> Deferred: + """Return a :class:`~twisted.internet.defer.Deferred` that fires when + a concurrency slot next becomes available (e.g. when + :meth:`record_done` is called or the limit is raised via + :meth:`set_concurrency`).""" + + def discard_slot_event(self, event: Deferred) -> None: + """Cancel a pending slot event returned by :meth:`slot_event`. + + Called by :class:`ThrottlingManager` when the wait ends without the + event firing (e.g. another scope's slot opened first). + """ + + def is_idle(self, now: float, max_idle: float) -> bool: + """Return whether this scope can be evicted from memory. + + A scope is idle when it has not been used for *max_idle* seconds and is + not currently in an active (future) backoff. + """ + + +# Safety timeout for acquire() while it waits, event-driven, for a concurrency +# slot to free up: a slot_event normally fires first (from record_done() or +# set_concurrency()); this only guards against a request that never reaches +# release() so the wait can never hang forever. +_SLOT_WAIT_TIMEOUT = 1.0 + class ThrottlingScopeManager: - """The default :setting:`THROTTLING_SCOPE_MANAGER` class.""" + """The default :setting:`THROTTLING_SCOPE_MANAGER` class. + + It implements a per-scope state machine covering delay, exponential + :ref:`backoff `, :ref:`rampup `, concurrency and + :ref:`quotas `: + + - A base :setting:`DOWNLOAD_DELAY`-style delay (``0`` by default, taken + from the scope ``"delay"`` config) is enforced between consecutive + requests for the scope. + + - On a backoff trigger (a :setting:`BACKOFF_HTTP_CODES` response or a + :setting:`BACKOFF_EXCEPTIONS` exception) the delay grows exponentially + by :setting:`BACKOFF_DELAY_FACTOR`, bounded by :setting:`BACKOFF_MIN_DELAY` + and :setting:`BACKOFF_MAX_DELAY`, with :setting:`BACKOFF_JITTER` applied. + A ``Retry-After`` / ``RateLimit-Reset`` delay is honored as a hard + minimum (capped at :setting:`BACKOFF_MAX_DELAY`). + + - After :setting:`BACKOFF_WINDOW` seconds without a new trigger, the delay + recovers one step at a time back towards the base delay. + + - When the scope is configured with a ``"concurrency"`` limit (or with + ``"rampup"``), no more than that many requests are allowed in flight at + once, never dropping below the ``"min_concurrency"`` floor. + + - When the scope sets ``"rampup": True``, throughput is increased every + :setting:`BACKOFF_WINDOW` that stays under :setting:`RAMPUP_BACKOFF_TARGET` + backoff triggers, first by lowering the delay and then by raising the + concurrency limit. + + - When the scope is configured with a ``"quota"``, no more than that much + quota is consumed per ``"window"`` (default: :setting:`THROTTLING_WINDOW`). + """ + + @classmethod + def from_crawler(cls, crawler: Crawler, config: dict[str, Any]) -> Self: + return cls(crawler, config) + + def __init__(self, crawler: Crawler, config: dict[str, Any]) -> None: + settings = crawler.settings + backoff: dict[str, Any] = config.get("backoff", {}) + self._id: ScopeID = config.get("id", "") + self._base_delay: float = float(config.get("delay", 0.0)) + self._randomize: bool = bool( + config.get("randomize_delay", settings.getbool("RANDOMIZE_DOWNLOAD_DELAY")) + ) + self._delay_factor: float = float( + backoff.get("delay_factor", settings.getfloat("BACKOFF_DELAY_FACTOR")) + ) + self._max_delay: float = float( + backoff.get("max_delay", settings.getfloat("BACKOFF_MAX_DELAY")) + ) + self._min_delay: float = float( + backoff.get("min_delay", settings.getfloat("BACKOFF_MIN_DELAY")) + ) + self._jitter: float | list[float] = backoff.get( + "jitter", settings.getfloat("BACKOFF_JITTER") + ) + self._window: float = settings.getfloat("BACKOFF_WINDOW") + self._min_concurrency: int = int(config.get("min_concurrency", 1)) + + # Rampup. + rampup = config.get("rampup") + self._rampup_enabled: bool = bool(rampup) + rampup_config: dict[str, Any] = rampup if isinstance(rampup, dict) else {} + self._rampup_target: tuple[float, float] = self._parse_target( + rampup_config.get("backoff_target", settings.get("RAMPUP_BACKOFF_TARGET")) + ) + self._rampup_delay_factor: float = float(rampup_config.get("delay_factor", 0.5)) + self._rampup_min_delay: float = float(rampup_config.get("min_delay", 0.0)) + + # Concurrency. ``None`` means no scope-level limit (the downloader slots + # enforce concurrency instead); a limit is only set when configured + # explicitly or implied by rampup. + configured_concurrency = config.get("concurrency") + if configured_concurrency is not None: + self._concurrency: int | None = int(configured_concurrency) + elif self._rampup_enabled: + self._concurrency = self._min_concurrency + else: + self._concurrency = None + + # Quota. + quota = config.get("quota") + self._quota: float | None = None if quota is None else float(quota) + self._quota_window: float = float( + config.get("window", settings.getfloat("THROTTLING_WINDOW")) + ) + + # State. + self._delay: float = self._base_delay + self._backoff_level: int = 0 + self._next_allowed_time: float | None = None + self._in_backoff_until: float | None = None + self._last_backoff_time: float | None = None + self._last_seen: float | None = None + self._active: int = 0 + self._slot_waiters: list[Deferred[None]] = [] + self._consumed: float = 0.0 + self._quota_window_start: float | None = None + self._rampup_window_start: float | None = None + self._rampup_backoffs: int = 0 + + @staticmethod + def _now(now: float | None) -> float: + return time.monotonic() if now is None else now + + @staticmethod + def _parse_target(value: Any) -> tuple[float, float]: + if isinstance(value, (list, tuple)): + return float(value[0]), float(value[1]) + return float(value), float(value) + + def _apply_jitter(self, value: float) -> float: + if isinstance(self._jitter, (list, tuple)): + low, high = self._jitter[0], self._jitter[1] + return value * (1 + random.uniform(low, high)) # noqa: S311 + if not self._jitter: + return value + return value * random.uniform(1 - self._jitter, 1 + self._jitter) # noqa: S311 + + def _effective_delay(self) -> float: + if self._backoff_level == 0 and self._randomize and self._delay > 0: + return random.uniform(0.5 * self._delay, 1.5 * self._delay) # noqa: S311 + return self._delay + + def _recover(self, now: float) -> None: + if self._backoff_level == 0 or self._last_backoff_time is None: + return + while self._backoff_level > 0 and now - self._last_backoff_time >= self._window: + self._backoff_level -= 1 + self._last_backoff_time += self._window + if self._backoff_level == 0: + self._delay = self._base_delay + self._in_backoff_until = None + self._last_backoff_time = None + break + self._delay = max(self._base_delay, self._delay / self._delay_factor) + + def _maybe_rampup(self, now: float) -> None: + """Increase throughput once per :setting:`BACKOFF_WINDOW` that stays + under :setting:`RAMPUP_BACKOFF_TARGET` backoff triggers.""" + if not self._rampup_enabled: + return + if self._rampup_window_start is None: + self._rampup_window_start = now + return + while now - self._rampup_window_start >= self._window: + self._rampup_window_start += self._window + if self._rampup_backoffs < self._rampup_target[0]: + self._rampup_step() + self._rampup_backoffs = 0 + + def _rampup_step(self) -> None: + # Backoff in progress: let it recover before probing again. + if self._backoff_level > 0: + return + if self._delay > self._rampup_min_delay: + self._delay = max( + self._rampup_min_delay, self._delay * self._rampup_delay_factor + ) + self._base_delay = min(self._base_delay, self._delay) + elif self._concurrency is not None: + self._concurrency += 1 + + def _maybe_reset_quota(self, now: float) -> None: + if self._quota is None: + return + if self._quota_window_start is None: + self._quota_window_start = now + return + while now - self._quota_window_start >= self._quota_window: + self._quota_window_start += self._quota_window + self._consumed = 0.0 + + def can_send(self, now: float | None = None, amount: float | None = None) -> float: + now = self._now(now) + self._recover(now) + self._maybe_rampup(now) + self._maybe_reset_quota(now) + waits = [0.0] + if self._in_backoff_until is not None: + waits.append(self._in_backoff_until - now) + if self._next_allowed_time is not None: + waits.append(self._next_allowed_time - now) + if self._quota is not None: + need = 0.0 if amount is None else float(amount) + # Block until the window resets only if some quota is already spent; + # a single oversized request is always allowed through. + if self._consumed > 0 and self._consumed + need > self._quota: + start = self._quota_window_start or now + waits.append(start + self._quota_window - now) + # Concurrency is enforced separately, via concurrency_blocked() and + # slot_event(), so acquire() can wait for a freed slot without polling. + return max(waits) + + def record_sent( + self, now: float | None = None, amount: float | None = None + ) -> None: + now = self._now(now) + self._last_seen = now + if self._in_backoff_until is not None and now >= self._in_backoff_until: + self._in_backoff_until = None + self._next_allowed_time = now + self._effective_delay() + self._active += 1 + if self._quota is not None and amount is not None: + self._maybe_reset_quota(now) + self._consumed += float(amount) + + def record_done(self, now: float | None = None) -> None: + if self._active > 0: + self._active -= 1 + self._fire_slot_waiters() + + def concurrency_blocked(self) -> bool: + return self._concurrency is not None and self._active >= self._concurrency + + def slot_event(self) -> Deferred[None]: + """Return a Deferred that fires when a concurrency slot next frees up + (via :meth:`record_done`) or the limit is raised (via + :meth:`set_concurrency`).""" + event: Deferred[None] = Deferred() + self._slot_waiters.append(event) + return event + + def discard_slot_event(self, event: Deferred[None]) -> None: + with contextlib.suppress(ValueError): + self._slot_waiters.remove(event) + + def _fire_slot_waiters(self) -> None: + waiters, self._slot_waiters = self._slot_waiters, [] + for event in waiters: + if not event.called: + event.callback(None) + + def record_backoff( + self, delay: float | None = None, now: float | None = None + ) -> None: + now = self._now(now) + self._last_seen = now + self._last_backoff_time = now + self._backoff_level += 1 + self._rampup_backoffs += 1 + if delay is not None: + hard = min(float(delay), self._max_delay) + self._in_backoff_until = now + hard + self._delay = min(max(self._delay, hard, self._min_delay), self._max_delay) + else: + grown = ( + self._delay * self._delay_factor if self._delay > 0 else self._min_delay + ) + grown = max(self._min_delay, grown) + grown = min(grown, self._max_delay) + self._delay = min(self._apply_jitter(grown), self._max_delay) + self._next_allowed_time = now + self._delay + + def reconcile_quota( + self, + consumed: float | None = None, + remaining: float | None = None, + now: float | None = None, + ) -> None: + if self._quota is None: + return + self._maybe_reset_quota(self._now(now)) + if remaining is not None: + self._consumed = max(0.0, self._quota - float(remaining)) + elif consumed is not None: + self._consumed = max(0.0, self._consumed + float(consumed)) + + def set_base_delay(self, delay: float) -> None: + if delay <= self._base_delay: + return + self._base_delay = delay + if self._backoff_level == 0: + self._delay = delay + + def set_concurrency(self, concurrency: int) -> None: + self._concurrency = max(self._min_concurrency, int(concurrency)) + self._fire_slot_waiters() + + def is_idle(self, now: float, max_idle: float) -> bool: + if self._in_backoff_until is not None and self._in_backoff_until > now: + return False + if self._active > 0: + return False + if self._last_seen is None: + return True + return (now - self._last_seen) > max_idle diff --git a/scrapy/utils/asyncio.py b/scrapy/utils/asyncio.py index 44604c0fe..71aecde9c 100644 --- a/scrapy/utils/asyncio.py +++ b/scrapy/utils/asyncio.py @@ -5,11 +5,11 @@ from __future__ import annotations import asyncio import logging import time -from collections.abc import AsyncIterator, Callable, Coroutine, Iterable +from collections.abc import AsyncIterator, Callable, Coroutine, Iterable, Sequence from typing import TYPE_CHECKING, Any, Concatenate, ParamSpec, TypeVar -from twisted.internet.defer import Deferred -from twisted.internet.task import LoopingCall +from twisted.internet.defer import Deferred, DeferredList +from twisted.internet.task import LoopingCall, deferLater from twisted.internet.threads import deferToThread from scrapy.utils.asyncgen import as_async_generator @@ -311,3 +311,88 @@ async def run_in_thread( from scrapy.utils.defer import maybe_deferred_to_future # noqa: PLC0415 return await maybe_deferred_to_future(deferToThread(func, *args, **kwargs)) + + +async def sleep(seconds: float) -> None: + """Sleep for *seconds*, working in asyncio-reactor, non-asyncio-reactor, + and reactorless modes. + + Uses :func:`asyncio.sleep` when asyncio is available, and + :func:`~twisted.internet.task.deferLater` otherwise. + + .. versionadded:: VERSION + """ + if is_asyncio_available(): + await asyncio.sleep(seconds) + else: + from twisted.internet import reactor + + # circular import + from scrapy.utils.defer import maybe_deferred_to_future # noqa: PLC0415 + + await maybe_deferred_to_future(deferLater(reactor, seconds, lambda: None)) + + +async def wait_for_first( + deferreds: Sequence[Deferred[Any]], + *, + timeout: float | None = None, +) -> tuple[set[Deferred[Any]], set[Deferred[Any]]]: + """Wait for the first of *deferreds* to fire, or until *timeout* seconds pass. + + Returns ``(done, pending)`` — two sets partitioning the input + :class:`~twisted.internet.defer.Deferred` objects — mirroring the API of + :func:`asyncio.wait`. + + Unfired deferreds in the ``pending`` set are neither cancelled nor + otherwise modified; the caller is responsible for any cleanup. + + Returns ``(set(), set())`` immediately when *deferreds* is empty. + + Works transparently in asyncio-reactor, non-asyncio-reactor, and + reactorless modes. + + .. versionadded:: VERSION + """ + if not deferreds: + return set(), set() + + if is_asyncio_available(): + # circular import + from scrapy.utils.defer import maybe_deferred_to_future # noqa: PLC0415 + + future_to_deferred = {maybe_deferred_to_future(d): d for d in deferreds} + done_futures, pending_futures = await asyncio.wait( + list(future_to_deferred), + timeout=timeout, + return_when=asyncio.FIRST_COMPLETED, + ) + return ( + {future_to_deferred[f] for f in done_futures}, + {future_to_deferred[f] for f in pending_futures}, + ) + + from twisted.internet import reactor + + # circular import + from scrapy.utils.defer import maybe_deferred_to_future # noqa: PLC0415 + + timeout_deferred = ( + deferLater(reactor, timeout, lambda: None) if timeout is not None else None + ) + waiter = DeferredList( + [*deferreds, *([] if timeout_deferred is None else [timeout_deferred])], + fireOnOneCallback=True, + fireOnOneErrback=True, + consumeErrors=True, + ) + try: + await maybe_deferred_to_future(waiter) + finally: + if timeout_deferred is not None and not timeout_deferred.called: + timeout_deferred.cancel() + + return ( + {d for d in deferreds if d.called}, + {d for d in deferreds if not d.called}, + ) diff --git a/tests/test_throttling.py b/tests/test_throttling.py new file mode 100644 index 000000000..0e0269c95 --- /dev/null +++ b/tests/test_throttling.py @@ -0,0 +1,567 @@ +from __future__ import annotations + +import logging + +import pytest + +from scrapy import signals +from scrapy.exceptions import DownloadTimeoutError +from scrapy.http import Request, Response +from scrapy.throttling import ThrottlingManager, ThrottlingScopeManager +from scrapy.utils.defer import deferred_from_coro, maybe_deferred_to_future +from scrapy.utils.test import get_crawler +from tests.spiders import SimpleSpider +from tests.utils.decorators import coroutine_test + + +def _manager(settings=None): + crawler = get_crawler(settings_dict=settings) + return ThrottlingManager.from_crawler(crawler) + + +def _scope_manager(settings=None, config=None): + crawler = get_crawler(settings_dict=settings) + return ThrottlingScopeManager.from_crawler(crawler, config or {"id": "example.com"}) + + +class _FakeRobotParser: + """A minimal robots.txt parser stub for :signal:`robots_parsed` tests. + + *delay* is returned by :meth:`crawl_delay`, unless it is an exception, in + which case it is raised to emulate a backend-specific failure. + """ + + def __init__(self, delay): + self._delay = delay + + def crawl_delay(self, useragent): + if isinstance(self._delay, Exception): + raise self._delay + return self._delay + + +def _response(status=200, headers=None, url="http://example.com", meta=None): + request = Request(url, meta=meta or {}) + return Response(url, status=status, headers=headers or {}, request=request) + + +class TestThrottlingManager: + @coroutine_test + async def test_get_scopes_returns_netloc(self): + manager = _manager() + assert ( + await manager.get_scopes(Request("http://example.com/a")) == "example.com" + ) + + @coroutine_test + async def test_get_scopes_cached(self): + manager = _manager() + request = Request("http://example.com/a") + first = await manager.get_scopes(request) + # A second call returns the cached value (same object identity for dicts, + # equal value for strings). + assert await manager.get_scopes(request) == first + + @coroutine_test + async def test_get_scopes_meta_string(self): + manager = _manager() + request = Request("http://example.com/a", meta={"throttling_scopes": "api"}) + assert await manager.get_scopes(request) == "api" + + @coroutine_test + async def test_get_scopes_meta_dict(self): + manager = _manager() + request = Request( + "http://example.com/a", meta={"throttling_scopes": {"api": 2.0}} + ) + assert await manager.get_scopes(request) == {"api": 2.0} + + @coroutine_test + async def test_get_initial_backoff_none(self): + manager = _manager() + assert await manager.get_initial_backoff() is None + + def test_scope_manager_class_in_config(self): + manager = _manager( + {"THROTTLING_SCOPES": {"example.com": {"manager": ThrottlingScopeManager}}} + ) + scope = manager._get_scope_manager("example.com") + assert isinstance(scope, ThrottlingScopeManager) + + def test_release_frees_concurrency(self): + manager = _manager({"THROTTLING_SCOPES": {"example.com": {"concurrency": 1}}}) + scope = manager._get_scope_manager("example.com") + request = Request("http://example.com") + scope.record_sent(now=0.0) + manager._reserved[request] = [(scope, None)] + assert scope.concurrency_blocked() is True + manager.release(request) + assert scope.concurrency_blocked() is False + # Releasing again is a no-op. + manager.release(request) + assert scope.concurrency_blocked() is False + + @coroutine_test + async def test_acquire_waits_for_freed_slot(self): + from scrapy.utils.asyncio import call_later # noqa: PLC0415 + + manager = _manager({"THROTTLING_SCOPES": {"example.com": {"concurrency": 1}}}) + r1 = Request("http://example.com/1") + r2 = Request("http://example.com/2") + # Drive acquire() the way the engine does, so it runs as a real task that + # can await the slot event under the asyncio reactor. + await maybe_deferred_to_future(deferred_from_coro(manager.acquire(r1))) + scope = manager._get_scope_manager("example.com") + assert scope.concurrency_blocked() is True + assert scope.concurrency_blocked() is True + # acquire(r2) must block until r1 frees the slot; release it on the next + # event loop tick so the event-driven wait wakes up. + call_later(0, manager.release, r1) + await maybe_deferred_to_future(deferred_from_coro(manager.acquire(r2))) + assert scope.concurrency_blocked() is True + assert r2 in manager._reserved + + def test_apply_backoff_reconciles_quota_without_backoff(self): + manager = _manager({"THROTTLING_SCOPES": {"cost": {"quota": 100.0}}}) + scope = manager._get_scope_manager("cost") + manager._apply_backoff({"cost": {"consumed": 5.0}}) + # Quota was reconciled but no backoff step was applied. + assert scope._consumed == pytest.approx(5.0) + assert scope._backoff_level == 0 + + def test_apply_backoff_delay_and_consumed(self): + manager = _manager({"THROTTLING_SCOPES": {"cost": {"quota": 100.0}}}) + scope = manager._get_scope_manager("cost") + manager._apply_backoff({"cost": {"delay": 5.0, "consumed": 2.0}}) + assert scope._consumed == pytest.approx(2.0) + assert scope._backoff_level == 1 + + @coroutine_test + async def test_response_backoff_non_backoff_code(self): + manager = _manager() + assert await manager.get_response_backoff(_response(status=200)) is None + + @coroutine_test + async def test_response_backoff_429_without_header(self): + manager = _manager() + assert ( + await manager.get_response_backoff(_response(status=429)) == "example.com" + ) + + @pytest.mark.parametrize( + ("header", "expected_delay"), + [ + ({"Retry-After": "7"}, 7.0), + ({"RateLimit-Reset": "12"}, 12.0), + ], + ids=["retry-after", "ratelimit-reset"], + ) + @coroutine_test + async def test_response_backoff_delay_header(self, header, expected_delay): + manager = _manager() + data = await manager.get_response_backoff(_response(status=429, headers=header)) + assert data == {"example.com": {"delay": expected_delay}} + + @coroutine_test + async def test_response_backoff_retry_after_http_date(self): + manager = _manager() + # A date far in the past yields no positive delay. + data = await manager.get_response_backoff( + _response( + status=503, + headers={"Retry-After": "Wed, 21 Oct 2015 07:28:00 GMT"}, + ) + ) + assert data == "example.com" + + @coroutine_test + async def test_response_backoff_max_of_both_headers(self): + manager = _manager() + data = await manager.get_response_backoff( + _response( + status=429, + headers={"Retry-After": "5", "RateLimit-Reset": "9"}, + ) + ) + assert data == {"example.com": {"delay": 9.0}} + + @coroutine_test + async def test_response_backoff_dont_track(self): + manager = _manager() + response = _response(status=429, meta={"throttling_dont_track": True}) + assert await manager.get_response_backoff(response) is None + + @pytest.mark.parametrize( + ("exception", "expected"), + [ + (DownloadTimeoutError(), "example.com"), + (ValueError(), None), + ], + ids=["tracked", "untracked"], + ) + @coroutine_test + async def test_exception_backoff(self, exception, expected): + manager = _manager() + request = Request("http://example.com") + assert await manager.get_exception_backoff(request, exception) == expected + + @coroutine_test + async def test_exception_backoff_dont_track(self): + manager = _manager() + request = Request("http://example.com", meta={"throttling_dont_track": True}) + assert ( + await manager.get_exception_backoff(request, DownloadTimeoutError()) is None + ) + + @pytest.mark.parametrize( + ("settings", "parser_delay", "expected_base_delay"), + [ + ({"ROBOTSTXT_OBEY": True, "RANDOMIZE_DOWNLOAD_DELAY": False}, 3.0, 3.0), + ({"ROBOTSTXT_OBEY": True, "THROTTLING_ROBOTSTXT_OBEY": False}, 3.0, 0.0), + ({"ROBOTSTXT_OBEY": True}, None, 0.0), + ({"ROBOTSTXT_OBEY": True}, ValueError(), 0.0), + ], + ids=["applies-delay", "obey-disabled", "no-delay", "backend-error"], + ) + def test_robots_parsed_signal(self, settings, parser_delay, expected_base_delay): + manager = _manager(settings) + manager.crawler.signals.send_catch_log( + signal=signals.robots_parsed, + robotparser=_FakeRobotParser(parser_delay), + request=Request("http://example.com/page"), + ) + scope = manager._get_scope_manager("example.com") + assert scope._base_delay == expected_base_delay + + def test_apply_robots_crawl_delay(self): + manager = _manager({"ROBOTSTXT_OBEY": True, "RANDOMIZE_DOWNLOAD_DELAY": False}) + manager.apply_robots_crawl_delay("example.com", 3.0) + scope = manager._get_scope_manager("example.com") + assert scope._base_delay == 3.0 + assert scope.can_send(now=0) == 0 # nothing sent yet + scope.record_sent(now=0) + assert scope.can_send(now=0) == pytest.approx(3.0) + + def test_apply_robots_crawl_delay_capped(self): + manager = _manager( + {"ROBOTSTXT_OBEY": True, "THROTTLING_ROBOTSTXT_MAX_DELAY": 2.0} + ) + manager.apply_robots_crawl_delay("example.com", 30.0) + assert manager._get_scope_manager("example.com")._base_delay == 2.0 + + def test_apply_robots_crawl_delay_disabled(self): + manager = _manager({"ROBOTSTXT_OBEY": True, "THROTTLING_ROBOTSTXT_OBEY": False}) + manager.apply_robots_crawl_delay("example.com", 3.0) + assert manager._get_scope_manager("example.com")._base_delay == 0.0 + + def test_apply_robots_crawl_delay_sets_concurrency(self): + manager = _manager({"ROBOTSTXT_OBEY": True}) + manager.apply_robots_crawl_delay("example.com", 3.0) + assert manager._get_scope_manager("example.com")._concurrency == 1 + + def test_apply_robots_crawl_delay_warns_on_conflict(self, caplog): + manager = _manager( + { + "ROBOTSTXT_OBEY": True, + "THROTTLING_SCOPES": {"example.com": {"concurrency": 8}}, + } + ) + with caplog.at_level(logging.WARNING, logger="scrapy.throttling"): + manager.apply_robots_crawl_delay("example.com", 3.0) + assert "Crawl-delay" in caplog.text + # The configured value takes precedence (crawl-delay not applied). + assert manager._get_scope_manager("example.com")._base_delay == 0.0 + + def test_apply_robots_crawl_delay_ignored(self, caplog): + manager = _manager( + { + "ROBOTSTXT_OBEY": True, + "THROTTLING_SCOPES": { + "example.com": {"concurrency": 8, "ignore_robots_txt": True} + }, + } + ) + with caplog.at_level(logging.WARNING, logger="scrapy.throttling"): + manager.apply_robots_crawl_delay("example.com", 3.0) + # No warning is logged and the crawl-delay is not applied. + assert "Crawl-delay" not in caplog.text + assert manager._get_scope_manager("example.com")._base_delay == 0.0 + + def test_scope_eviction(self): + manager = _manager({"THROTTLING_SCOPE_MAX_IDLE": 100.0}) + scope = manager._get_scope_manager("example.com") + scope.record_sent(now=0.0) + scope.record_done(now=0.0) + # Not idle yet. + manager._last_eviction = None + manager._maybe_evict(now=50.0) + assert "example.com" in manager._scope_managers + # Idle past the threshold. + manager._last_eviction = None + manager._maybe_evict(now=201.0) + assert "example.com" not in manager._scope_managers + + def test_scope_eviction_skips_active_backoff(self): + manager = _manager( + {"THROTTLING_SCOPE_MAX_IDLE": 100.0, "BACKOFF_MAX_DELAY": 100_000.0} + ) + scope = manager._get_scope_manager("example.com") + scope.record_backoff(delay=10_000.0, now=0.0) + manager._last_eviction = None + manager._maybe_evict(now=5_000.0) + # Still in backoff (in_backoff_until far in the future), so not evicted + # even though it has been idle for longer than THROTTLING_SCOPE_MAX_IDLE. + assert "example.com" in manager._scope_managers + + +class TestThrottlingScopeManager: + def test_no_delay_by_default(self): + scope = _scope_manager() + scope.record_sent(now=0.0) + assert scope.can_send(now=0.0) == 0 + + def test_base_delay_enforced(self): + scope = _scope_manager( + {"RANDOMIZE_DOWNLOAD_DELAY": False}, {"id": "x", "delay": 2.0} + ) + scope.record_sent(now=10.0) + assert scope.can_send(now=10.0) == pytest.approx(2.0) + assert scope.can_send(now=11.0) == pytest.approx(1.0) + assert scope.can_send(now=12.0) == 0 + + def test_exponential_backoff(self): + scope = _scope_manager( + { + "BACKOFF_MIN_DELAY": 1.0, + "BACKOFF_DELAY_FACTOR": 2.0, + "BACKOFF_JITTER": 0, + }, + ) + scope.record_backoff(now=0.0) + assert scope._delay == pytest.approx(1.0) + scope.record_backoff(now=0.0) + assert scope._delay == pytest.approx(2.0) + scope.record_backoff(now=0.0) + assert scope._delay == pytest.approx(4.0) + + def test_backoff_cap(self): + scope = _scope_manager( + { + "BACKOFF_MIN_DELAY": 1.0, + "BACKOFF_DELAY_FACTOR": 10.0, + "BACKOFF_MAX_DELAY": 5.0, + "BACKOFF_JITTER": 0, + }, + ) + for _ in range(5): + scope.record_backoff(now=0.0) + assert scope._delay == pytest.approx(5.0) + + @pytest.mark.parametrize( + ("max_delay", "backoff_delay", "expected"), + [ + (100.0, 20.0, 20.0), + (10.0, 999.0, 10.0), + ], + ids=["within-cap", "capped"], + ) + def test_retry_after_delay(self, max_delay, backoff_delay, expected): + scope = _scope_manager({"BACKOFF_MAX_DELAY": max_delay}) + scope.record_backoff(delay=backoff_delay, now=0.0) + assert scope.can_send(now=0.0) == pytest.approx(expected) + + def test_recovery_after_window(self): + scope = _scope_manager( + { + "BACKOFF_MIN_DELAY": 1.0, + "BACKOFF_DELAY_FACTOR": 2.0, + "BACKOFF_WINDOW": 60.0, + "BACKOFF_JITTER": 0, + }, + ) + scope.record_backoff(now=0.0) + scope.record_backoff(now=0.0) + assert scope._delay == pytest.approx(2.0) + # One window passes with no new backoff -> one step down. + scope.can_send(now=60.0) + assert scope._backoff_level == 1 + assert scope._delay == pytest.approx(1.0) + # Another window -> back to base (0). + scope.can_send(now=120.0) + assert scope._backoff_level == 0 + assert scope._delay == pytest.approx(0.0) + + def test_per_scope_backoff_override(self): + scope = _scope_manager( + {"BACKOFF_MIN_DELAY": 1.0, "BACKOFF_DELAY_FACTOR": 2.0}, + { + "id": "x", + "backoff": {"min_delay": 5.0, "delay_factor": 3.0, "jitter": 0}, + }, + ) + scope.record_backoff(now=0.0) + assert scope._delay == pytest.approx(5.0) + scope.record_backoff(now=0.0) + assert scope._delay == pytest.approx(15.0) + + def test_set_base_delay_raises_only(self): + scope = _scope_manager( + {"RANDOMIZE_DOWNLOAD_DELAY": False}, {"id": "x", "delay": 5.0} + ) + scope.set_base_delay(2.0) # lower -> ignored + assert scope._base_delay == 5.0 + scope.set_base_delay(8.0) # higher -> applied + assert scope._base_delay == 8.0 + assert scope._delay == 8.0 + + def test_min_delay_first_step(self): + scope = _scope_manager( + {"BACKOFF_MIN_DELAY": 3.0, "BACKOFF_DELAY_FACTOR": 2.0, "BACKOFF_JITTER": 0} + ) + scope.record_backoff(now=0.0) + assert scope._delay == pytest.approx(3.0) + + def test_no_scope_concurrency_limit_by_default(self): + scope = _scope_manager() + assert scope._concurrency is None + for _ in range(100): + scope.record_sent(now=0.0) + assert scope.can_send(now=0.0) == 0 + assert scope.concurrency_blocked() is False + + def test_concurrency_limit(self): + scope = _scope_manager(config={"id": "x", "concurrency": 2}) + scope.record_sent(now=0.0) + # Concurrency is enforced via concurrency_blocked(), not can_send(). + assert scope.can_send(now=0.0) == 0 + assert scope.concurrency_blocked() is False + scope.record_sent(now=0.0) + # Two in flight, limit reached -> blocked. + assert scope.can_send(now=0.0) == 0 + assert scope.concurrency_blocked() is True + scope.record_done(now=0.0) + assert scope.concurrency_blocked() is False + + def test_record_done_fires_slot_event(self): + scope = _scope_manager(config={"id": "x", "concurrency": 1}) + scope.record_sent(now=0.0) + event = scope.slot_event() + assert not event.called + scope.record_done(now=0.0) + assert event.called + + def test_set_concurrency_fires_slot_event(self): + scope = _scope_manager(config={"id": "x", "concurrency": 1}) + scope.record_sent(now=0.0) + event = scope.slot_event() + assert not event.called + scope.set_concurrency(5) + assert event.called + + def test_discard_slot_event(self): + scope = _scope_manager(config={"id": "x", "concurrency": 1}) + event = scope.slot_event() + scope.discard_slot_event(event) + scope.discard_slot_event(event) # idempotent + scope.record_sent(now=0.0) + scope.record_done(now=0.0) + assert not event.called + + def test_set_concurrency_respects_min(self): + scope = _scope_manager(config={"id": "x", "min_concurrency": 3}) + scope.set_concurrency(1) + assert scope._concurrency == 3 + scope.set_concurrency(5) + assert scope._concurrency == 5 + + def test_quota_blocks_when_exhausted(self): + scope = _scope_manager(config={"id": "x", "quota": 10.0, "window": 60.0}) + scope.record_sent(now=0.0, amount=6.0) + assert scope.can_send(now=0.0, amount=3.0) == 0 # 9 <= 10 + scope.record_sent(now=0.0, amount=3.0) + # 9 spent; a 3.0 request would exceed the quota -> wait for the window. + assert scope.can_send(now=0.0, amount=3.0) == pytest.approx(60.0) + # The window resets and quota is available again. + assert scope.can_send(now=60.0, amount=3.0) == 0 + + def test_quota_allows_oversized_request(self): + scope = _scope_manager(config={"id": "x", "quota": 10.0}) + # A single request larger than the whole quota is still allowed. + assert scope.can_send(now=0.0, amount=999.0) == 0 + + def test_quota_reconcile_consumed_delta(self): + scope = _scope_manager(config={"id": "x", "quota": 10.0}) + scope.record_sent(now=0.0, amount=2.0) + assert scope._consumed == pytest.approx(2.0) + # The response reports it actually consumed 0.5 more than estimated. + scope.reconcile_quota(consumed=0.5, now=0.0) + assert scope._consumed == pytest.approx(2.5) + + def test_quota_reconcile_remaining(self): + scope = _scope_manager(config={"id": "x", "quota": 10.0}) + scope.record_sent(now=0.0, amount=2.0) + scope.reconcile_quota(remaining=3.0, now=0.0) + assert scope._consumed == pytest.approx(7.0) + + def test_rampup_lowers_delay_when_quiet(self): + scope = _scope_manager( + {"BACKOFF_WINDOW": 10.0, "RANDOMIZE_DOWNLOAD_DELAY": False}, + { + "id": "x", + "delay": 4.0, + "rampup": {"delay_factor": 0.5, "min_delay": 0.5}, + }, + ) + scope.can_send(now=0.0) # start the rampup window + # A quiet window (no backoff) lowers the delay. + scope.can_send(now=10.0) + assert scope._delay == pytest.approx(2.0) + scope.can_send(now=20.0) + assert scope._delay == pytest.approx(1.0) + + def test_rampup_raises_concurrency_at_min_delay(self): + scope = _scope_manager( + {"BACKOFF_WINDOW": 10.0}, + {"id": "x", "delay": 0.0, "rampup": True, "min_concurrency": 1}, + ) + assert scope._concurrency == 1 + scope.can_send(now=0.0) + scope.can_send(now=10.0) + assert scope._concurrency == 2 + + def test_rampup_holds_when_target_met(self): + scope = _scope_manager( + {"BACKOFF_WINDOW": 10.0, "RAMPUP_BACKOFF_TARGET": 1}, + {"id": "x", "delay": 0.0, "rampup": True, "min_concurrency": 1}, + ) + scope.can_send(now=0.0) + scope.record_backoff(now=1.0) # one trigger == target -> hold, do not probe + scope.can_send(now=10.0) + assert scope._concurrency == 1 + + +class TestThrottlingIntegration: + @coroutine_test + async def test_backoff_recorded_on_429(self, mockserver): + crawler = get_crawler(SimpleSpider, {"RETRY_ENABLED": False}) + await crawler.crawl_async( + mockserver.url("/status?n=429"), mockserver=mockserver + ) + throttler = crawler.throttler + assert throttler is not None + managers = throttler._scope_managers + assert managers, "no throttling scope was created" + assert any(manager._backoff_level >= 1 for manager in managers.values()) + + @coroutine_test + async def test_no_backoff_on_200(self, mockserver): + crawler = get_crawler(SimpleSpider) + await crawler.crawl_async( + mockserver.url("/status?n=200"), mockserver=mockserver + ) + throttler = crawler.throttler + assert throttler is not None + assert all( + manager._backoff_level == 0 + for manager in throttler._scope_managers.values() + )