From 074f9e567dfa86bf647b695b4725e48f59366a66 Mon Sep 17 00:00:00 2001 From: Adrian Chaves Date: Wed, 1 Jul 2026 10:19:27 +0200 Subject: [PATCH] WIP --- docs/conf.py | 1 - docs/topics/downloader-middleware.rst | 22 ++ docs/topics/throttling.rst | 87 ++++-- scrapy/core/engine.py | 15 +- scrapy/downloadermiddlewares/backoff.py | 169 ++++++++++ scrapy/settings/default_settings.py | 1 + scrapy/throttling.py | 397 +++++++----------------- 7 files changed, 357 insertions(+), 335 deletions(-) create mode 100644 scrapy/downloadermiddlewares/backoff.py diff --git a/docs/conf.py b/docs/conf.py index 2fc7afe8c..7770e550b 100644 --- a/docs/conf.py +++ b/docs/conf.py @@ -144,7 +144,6 @@ coverage_ignore_pyobjects = [ # -- Options for the autodoc extension ---------------------------------------- autodoc_type_aliases = { - "BackoffData": "BackoffData", "RequestScopes": "RequestScopes", } diff --git a/docs/topics/downloader-middleware.rst b/docs/topics/downloader-middleware.rst index b2259eacb..83f28cf2a 100644 --- a/docs/topics/downloader-middleware.rst +++ b/docs/topics/downloader-middleware.rst @@ -169,6 +169,28 @@ middleware, see the :ref:`downloader middleware usage guide For a list of the components enabled by default (and their orders) see the :setting:`DOWNLOADER_MIDDLEWARES_BASE` setting. +.. _backoff-mw: + +BackoffMiddleware +----------------- + +.. module:: scrapy.downloadermiddlewares.backoff + :synopsis: Backoff Downloader Middleware + +.. class:: BackoffMiddleware + + This middleware feeds download outcomes into :ref:`throttling `: + for every response matching :setting:`BACKOFF_HTTP_CODES` and every download + exception matching :setting:`BACKOFF_EXCEPTIONS`, it makes the request's + :ref:`throttling scopes ` :ref:`back off `. + + It runs below + :class:`~scrapy.downloadermiddlewares.retry.RetryMiddleware` so that it + observes rate-limiting responses (``429``, ``503``, …) before they are + turned into retries. + + See :ref:`throttling` for details. + .. _cookies-mw: CookiesMiddleware diff --git a/docs/topics/throttling.rst b/docs/topics/throttling.rst index 1c82a7c93..09b8e4113 100644 --- a/docs/topics/throttling.rst +++ b/docs/topics/throttling.rst @@ -165,6 +165,14 @@ full speed gradually once it recovers. To keep a scope hovering around a target rate instead of repeatedly probing and backing off, enable :ref:`rampup `. +Backoff triggers are detected by the +:class:`~scrapy.downloadermiddlewares.backoff.BackoffMiddleware`, a built-in +:ref:`downloader middleware ` enabled by default. +Any component can also trigger backoff programmatically for arbitrary scopes — +e.g. based on the response body of a specific site — through +:meth:`crawler.throttler.back_off() +`. + .. _per-scope-backoff: Per-scope backoff configuration @@ -793,32 +801,43 @@ to: target_domain = urlparse(target_url).netloc return add_scope(scopes, target_domain) - - Can differentiate between exhaustion of the target website and - exhaustion of the API itself. For example: +- Add a :ref:`downloader middleware ` that + differentiates between exhaustion of the target website and exhaustion of + the API itself. The API returns ``200`` even when the target website + rate-limits it, reporting the upstream status in a header; the middleware + backs off the **target-website** scope (not the API scope) in that case, + reusing the scope's own :meth:`~scrapy.throttling.ThrottlingScopeManagerProtocol.triggers_backoff_for_status` + check: - .. code-block:: python + .. code-block:: python - from scrapy.throttling import ThrottlingManager - from scrapy.utils.httpobj import urlparse_cached + from scrapy.throttling import iter_scopes + from scrapy.utils.httpobj import urlparse_cached - class MyThrottlingManager(ThrottlingManager): - async def get_response_backoff(self, response): - if ( - urlparse_cached(response.request).netloc != "api.toscrape.com" - or response.status != 200 - ): - return await super().get_response_backoff(response) - upstream_status_code = int( + class UpstreamBackoffMiddleware: + def __init__(self, crawler): + self.throttler = crawler.throttler + + @classmethod + def from_crawler(cls, crawler): + return cls(crawler) + + def process_response(self, request, response, spider): + if urlparse_cached(request).netloc == "api.toscrape.com": + upstream_status = int( response.headers.get("X-Upstream-Status-Code", b"200") ) - upstream_response = response.__class__( - response.url, - status=upstream_status_code, - headers=response.headers, - body=response.body, - ) - return await super().get_response_backoff(upstream_response) + scopes = [ + scope + for scope in iter_scopes(self.throttler.get_resolved_scopes(request)) + if scope != "api.toscrape.com" + and self.throttler.get_scope_manager(scope).triggers_backoff_for_status( + upstream_status + ) + ] + self.throttler.back_off(scopes) + return response .. _cost-smoothing-throttling: @@ -851,23 +870,27 @@ window (:setting:`THROTTLING_WINDOW`). You can use :ref:`throttling quotas return scopes return add_scope(scopes, "cost", estimate_request_cost(request)) - - Reconciles the estimated cost with the actual cost reported by the - response, so that the quota tracks real spending: +- Add a :ref:`downloader middleware ` that + reconciles the estimated cost with the actual cost reported by the + response, so that the quota tracks real spending: - .. code-block:: python + .. code-block:: python - from scrapy.throttling import ThrottlingManager, update_scope_backoff + class CostReconcileMiddleware: + def __init__(self, crawler): + self.throttler = crawler.throttler + @classmethod + def from_crawler(cls, crawler): + return cls(crawler) - class MyThrottlingManager(ThrottlingManager): - async def get_response_backoff(self, response): - backoff = await super().get_response_backoff(response) - if response.headers.get("X-Actual-Cost") is None: - return backoff - estimated = estimate_request_cost(response.request) + def process_response(self, request, response, spider): + if response.headers.get("X-Actual-Cost") is not None: + estimated = estimate_request_cost(request) actual = float(response.headers[b"X-Actual-Cost"]) # Report the difference between actual and estimated cost. - return update_scope_backoff(backoff, "cost", consumed=actual - estimated) + self.throttler.reconcile_quota("cost", consumed=actual - estimated) + return response - Use the :setting:`THROTTLING_SCOPES` setting to set a maximum cost per time window: @@ -1101,7 +1124,6 @@ API :member-order: bysource .. autoclass:: scrapy.throttling.ThrottlingManager - :members: get_response_delay .. autoclass:: scrapy.throttling.ThrottlingScopeManagerProtocol :members: @@ -1119,4 +1141,3 @@ API .. autofunction:: scrapy.throttling.scope_cache .. autofunction:: scrapy.throttling.add_scope -.. autofunction:: scrapy.throttling.update_scope_backoff diff --git a/scrapy/core/engine.py b/scrapy/core/engine.py index 75a70ce9a..b18ebf0d0 100644 --- a/scrapy/core/engine.py +++ b/scrapy/core/engine.py @@ -627,15 +627,11 @@ class ExecutionEngine: assert throttler is not None try: yield self._acquire_throttling(request) - try: - result: Response | Request - if self._downloader_fetch_needs_spider: - result = yield self.downloader.fetch(request, self.spider) - else: - result = yield self.downloader.fetch(request) - except Exception as exc: - yield deferred_from_coro(throttler.process_exception(request, exc)) - raise + result: Response | Request + if self._downloader_fetch_needs_spider: + result = yield self.downloader.fetch(request, self.spider) + else: + result = yield self.downloader.fetch(request) if not isinstance(result, (Response, Request)): raise TypeError( f"Incorrect type: expected Response or Request, got {type(result)}: {result!r}" @@ -643,7 +639,6 @@ class ExecutionEngine: if isinstance(result, Response): if result.request is None: result.request = request - yield deferred_from_coro(throttler.process_response(result)) logkws = self.logformatter.crawled(result.request, result, self.spider) if logkws is not None: logger.log( diff --git a/scrapy/downloadermiddlewares/backoff.py b/scrapy/downloadermiddlewares/backoff.py new file mode 100644 index 000000000..173ef1e62 --- /dev/null +++ b/scrapy/downloadermiddlewares/backoff.py @@ -0,0 +1,169 @@ +from __future__ import annotations + +import datetime as dt +import logging +from email.utils import parsedate_to_datetime +from typing import TYPE_CHECKING + +from scrapy.throttling import iter_scopes +from scrapy.utils.decorators import _warn_spider_arg +from scrapy.utils.misc import load_object + +if TYPE_CHECKING: + # typing.Self requires Python 3.11 + from typing_extensions import Self + + import scrapy + from scrapy.crawler import Crawler + from scrapy.http import Request, Response + from scrapy.throttling import ThrottlingManagerProtocol + + +logger = logging.getLogger(__name__) + + +def _parse_retry_after(response: Response) -> float | None: + raw = response.headers.get("Retry-After") + if not raw: + return None + try: + value = raw.decode("utf-8").strip() + except UnicodeDecodeError: + return None + if value.isdigit(): + return float(value) # seconds + try: + date = parsedate_to_datetime(value) + except (TypeError, ValueError, OverflowError): + return None + if date.tzinfo is None: + date = date.replace(tzinfo=dt.timezone.utc) + now = dt.datetime.now(dt.timezone.utc) + seconds_to_wait = (date - now).total_seconds() + # Keep sub-second precision (a date less than a second away must not be + # truncated to 0 and dropped); a past or present date yields no delay. + return max(0.0, seconds_to_wait) or None + + +def _parse_ratelimit_reset(response: Response) -> float | None: + raw = response.headers.get("RateLimit-Reset") + if not raw: + return None + try: + value = raw.decode("utf-8").strip() + except UnicodeDecodeError: + return None + try: + return float(value) + except ValueError: + return None + + +class BackoffMiddleware: + """Downloader middleware that drives :ref:`backoff ` from download + outcomes. + + It observes every response and download exception and, for those matching + :setting:`BACKOFF_HTTP_CODES` or :setting:`BACKOFF_EXCEPTIONS` (globally or + per :setting:`THROTTLING_SCOPES` scope), tells the :ref:`throttling manager + ` to back off the request's scopes through its + :meth:`~scrapy.throttling.ThrottlingManagerProtocol.back_off` API. + + It sits below :class:`~scrapy.downloadermiddlewares.retry.RetryMiddleware` + in :setting:`DOWNLOADER_MIDDLEWARES_BASE` so it sees rate-limiting responses + (429, 503, …) before the retry middleware turns them into new requests. + """ + + def __init__(self, crawler: Crawler): + # Throttling is a core, always-on subsystem: THROTTLING_MANAGER has a + # non-None default and is instantiated before the downloader is built, + # so crawler.throttler is always set here (the engine likewise asserts + # it in its download path). + assert crawler.throttler is not None + self._throttler: ThrottlingManagerProtocol = crawler.throttler + settings = crawler.settings + # Union of the global backoff triggers and every per-scope override: a + # response status (or exception type) outside it cannot trigger backoff + # for any scope, so the scopes of such a request need not be resolved. + # Each scope still makes the final decision via its scope manager's + # triggers_backoff_* methods (which read the per-scope overrides). + self._http_codes: set[int] = { + int(code) for code in settings.getlist("BACKOFF_HTTP_CODES") + } + self._exceptions: tuple[type[BaseException], ...] = tuple( + load_object(exc) if isinstance(exc, str) else exc + for exc in settings.getlist("BACKOFF_EXCEPTIONS") + ) + for scope_config in settings.getdict("THROTTLING_SCOPES").values(): + backoff = scope_config.get("backoff") or {} + if "http_codes" in backoff: + self._http_codes.update(int(code) for code in backoff["http_codes"]) + if "exceptions" in backoff: + self._exceptions += tuple( + load_object(exc) if isinstance(exc, str) else exc + for exc in backoff["exceptions"] + ) + + @classmethod + def from_crawler(cls, crawler: Crawler) -> Self: + return cls(crawler) + + @_warn_spider_arg + def process_response( + self, + request: Request, + response: Response, + spider: scrapy.Spider | None = None, + ) -> Response: + if ( + response.status not in self._http_codes + or "cached" in response.flags + or request.meta.get("throttling_dont_track") + ): + return response + matched = [ + scope + for scope in iter_scopes(self._throttler.get_resolved_scopes(request)) + if self._throttler.get_scope_manager(scope).triggers_backoff_for_status( + response.status + ) + ] + if matched: + self._throttler.back_off(matched, delay=self._response_delay(response)) + return response + + @_warn_spider_arg + def process_exception( + self, + request: Request, + exception: Exception, + spider: scrapy.Spider | None = None, + ) -> None: + if request.meta.get("throttling_dont_track") or not isinstance( + exception, self._exceptions + ): + return + matched = [ + scope + for scope in iter_scopes(self._throttler.get_resolved_scopes(request)) + if self._throttler.get_scope_manager(scope).triggers_backoff_for_exception( + exception + ) + ] + if matched: + self._throttler.back_off(matched) + return + + @staticmethod + def _response_delay(response: Response) -> float | None: + """Return the hard minimum delay requested by *response* through a + ``Retry-After`` or ``RateLimit-Reset`` header, or ``None``.""" + delays = [ + delay + for delay in ( + _parse_retry_after(response), + _parse_ratelimit_reset(response), + ) + if delay is not None + ] + return max(delays) if delays else None diff --git a/scrapy/settings/default_settings.py b/scrapy/settings/default_settings.py index 6dd6160b6..f6e677e8b 100644 --- a/scrapy/settings/default_settings.py +++ b/scrapy/settings/default_settings.py @@ -359,6 +359,7 @@ DOWNLOADER_MIDDLEWARES_BASE = { "scrapy.downloadermiddlewares.redirect.MetaRefreshMiddleware": 580, "scrapy.downloadermiddlewares.httpcompression.HttpCompressionMiddleware": 590, "scrapy.downloadermiddlewares.redirect.RedirectMiddleware": 600, + "scrapy.downloadermiddlewares.backoff.BackoffMiddleware": 630, "scrapy.downloadermiddlewares.cookies.CookiesMiddleware": 700, "scrapy.downloadermiddlewares.httpproxy.HttpProxyMiddleware": 750, "scrapy.downloadermiddlewares.stats.DownloaderStats": 850, diff --git a/scrapy/throttling.py b/scrapy/throttling.py index dcdfeb97d..5b4079746 100644 --- a/scrapy/throttling.py +++ b/scrapy/throttling.py @@ -1,7 +1,6 @@ from __future__ import annotations import contextlib -import datetime as dt import ipaddress import logging import random @@ -10,13 +9,12 @@ import time import warnings from collections import OrderedDict from collections.abc import Awaitable, Callable, Iterable -from email.utils import parsedate_to_datetime from functools import wraps 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 typing_extensions import Self from scrapy import signals from scrapy.exceptions import ScrapyDeprecationWarning @@ -27,56 +25,13 @@ from scrapy.utils.misc import build_from_crawler, load_object if TYPE_CHECKING: from scrapy.crawler import Crawler - from scrapy.http import Request, Response + from scrapy.http import Request from scrapy.settings import BaseSettings logger = logging.getLogger(__name__) -def _parse_retry_after(response: Response) -> float | None: - raw = response.headers.get("Retry-After") - if not raw: - return None - try: - value = raw.decode("utf-8").strip() - except UnicodeDecodeError: - return None - if value.isdigit(): - return float(value) # seconds - try: - date = parsedate_to_datetime(value) - except (TypeError, ValueError, OverflowError): - return None - if date.tzinfo is None: - date = date.replace(tzinfo=dt.timezone.utc) - now = dt.datetime.now(dt.timezone.utc) - seconds_to_wait = (date - now).total_seconds() - # Keep sub-second precision (a date less than a second away must not be - # truncated to 0 and dropped); a past or present date yields no delay. - return max(0.0, seconds_to_wait) or None - - -def _parse_ratelimit_reset(response: Response) -> float | None: - raw = response.headers.get("RateLimit-Reset") - if not raw: - return None - try: - value = raw.decode("utf-8").strip() - except UnicodeDecodeError: - return None - try: - return float(value) - except ValueError: - return None - - -class BackoffScopeData(TypedDict): - delay: NotRequired[float] - consumed: NotRequired[float] - remaining: NotRequired[float] - - class BackoffConfig(TypedDict, total=False): """Per-scope override of the backoff settings. @@ -142,7 +97,6 @@ class ThrottlingScopeConfig(TypedDict, total=False): ScopeID = str -BackoffData = None | ScopeID | Iterable[ScopeID] | dict[ScopeID, BackoffScopeData] RequestScopes = None | ScopeID | Iterable[ScopeID] | dict[ScopeID, float | None] @@ -285,36 +239,6 @@ def add_scope( return scopes -def update_scope_backoff( - backoff: BackoffData, - scope: ScopeID, - /, - *, - delay: float | None = None, - consumed: float | None = None, -) -> BackoffData: - """Add *scope* to *backoff* or update its existing entry with the given - parameters. - - This is a utility function to help extending the output of - :meth:`~ThrottlingManagerProtocol.get_initial_backoff`, - :meth:`~ThrottlingManagerProtocol.get_response_backoff` or - :meth:`~ThrottlingManagerProtocol.get_exception_backoff`, e.g. in - :class:`ThrottlingManager` subclasses. - """ - 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 - - class ThrottlingManagerProtocol(Protocol): """A protocol for :setting:`THROTTLING_MANAGER` :ref:`components `.""" @@ -328,64 +252,17 @@ class ThrottlingManagerProtocol(Protocol): keys and :ref:`throttling quotas ` as values. """ - async def get_initial_backoff(self) -> BackoffData: - """Return the initial throttling data. + def get_resolved_scopes(self, request: Request) -> RequestScopes: + """Return the :ref:`throttling scopes ` under which + *request* was (or will be) sent, without re-resolving them. - This method is called before the first request is sent, and it should - be used to provide an initial throttling state, to be used before it is - updated with later calls to :meth:`get_response_backoff` and - :meth:`get_exception_backoff`. - - **Return values:** - - You may return any of the following: - - - ``None``: no throttling data to report. - - - A string: a single scope name, indicating that the scope is - currently exhausted. - - - An iterable of strings: multiple scope names, indicating that - those scopes are currently exhausted. - - - A dict with scope names as keys and dict values. Dict values - support the following keys (matching :class:`BackoffScopeData`): - - - ``"delay"``: a float indicating how many seconds to wait before - sending another request for the scope. - - - ``"remaining"``: a float indicating the remaining - :ref:`throttling quota ` of the scope. - - - ``"consumed"``: a float indicating how much of the scope's - quota has already been consumed. - - An empty dict marks the scope as currently exhausted. - - For example: - - .. code-block:: python - - return { - "scope1": {"delay": 5.0}, - "scope2": {}, - "scope3": {"remaining": 42.0}, - } - """ - - async def get_response_backoff(self, response: Response) -> BackoffData: - """Return a throttling data update based on *response*. - - It supports the same return values as :meth:`get_initial_backoff`. - """ - - async def get_exception_backoff( - self, request: Request, exception: Exception - ) -> BackoffData: - """Return a throttling data update based on *exception* and the - *request* that caused it. - - It supports the same return values as :meth:`get_initial_backoff`. + This is the synchronous counterpart of :meth:`get_scopes`: it returns + the scopes resolved earlier (e.g. at enqueue or :meth:`acquire` time) + and persisted on ``request.meta``, falling back to a best-effort + synchronous resolution only if none were persisted. Use it, rather than + :meth:`get_scopes`, to attribute a response or exception to the very + scopes the request was sent under — e.g. from a downloader middleware or + a spider callback that wants to :meth:`back_off` based on the response. """ async def acquire(self, request: Request) -> None: @@ -467,17 +344,62 @@ class ThrottlingManagerProtocol(Protocol): that share its scopes. """ + def back_off( + self, + scopes: RequestScopes, + *, + delay: float | None = None, + cap: bool = True, + ) -> None: + """Register a :ref:`backoff ` trigger for each of *scopes*. + + This is the general-purpose way to make a scope slow down, available to + any component through :attr:`crawler.throttler + `. The built-in :class:`backoff + middleware ` + calls it for :setting:`BACKOFF_HTTP_CODES` responses and + :setting:`BACKOFF_EXCEPTIONS` exceptions, but a downloader middleware or + spider callback can call it too (e.g. to back off based on the response + body of a specific site). + + *scopes* accepts the same shapes as the output of :meth:`get_scopes` + (typically the result of :meth:`get_resolved_scopes` for a request). + + When *delay* is ``None`` an exponential backoff step is applied; when + given, *delay* is a hard minimum delay in seconds (e.g. from a + :ref:`Retry-After ` header). *cap* limits *delay* to + :setting:`BACKOFF_MAX_DELAY`; set it to ``False`` for trusted, + programmatic delays. + """ + def delay_scope(self, scope_id: str, delay: float) -> None: """Hold back every request of the scope identified by *scope_id* for at least *delay* seconds, counted as a :ref:`backoff ` trigger for the scope. - This is the programmatic equivalent of a :ref:`Retry-After - ` response header, available to any component through - :attr:`crawler.throttler `. Unlike - those headers, *delay* is **not** capped at - :setting:`BACKOFF_MAX_DELAY`: that cap guards against untrusted input, - whereas a ``delay_scope`` call is trusted. + This is shorthand for :meth:`back_off(scope_id, delay=delay, + cap=False) `: the programmatic equivalent of a + :ref:`Retry-After ` response header. Unlike those headers, + *delay* is **not** capped at :setting:`BACKOFF_MAX_DELAY`: that cap + guards against untrusted input, whereas a ``delay_scope`` call is + trusted. + """ + + def reconcile_quota( + self, + scopes: RequestScopes, + *, + consumed: float | None = None, + remaining: float | None = None, + ) -> None: + """Reconcile the :ref:`throttling quota ` of each of + *scopes* with an actually *consumed* amount (a delta to add) or a + *remaining* amount (an absolute value), correcting the estimate used + when requests were sent. + + Like :meth:`back_off`, this is meant to be called from a downloader + middleware or spider callback that learns the real quota cost of a + request from its response. """ def get_scope_delay(self, scope_id: str) -> float: @@ -495,12 +417,6 @@ class ThrottlingManagerProtocol(Protocol): """Return the :class:`ThrottlingScopeManagerProtocol` instance handling the scope identified by *scope_id*, creating it if necessary.""" - 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*.""" - _GetScopesMethod = TypeVar( "_GetScopesMethod", bound=Callable[..., Awaitable[RequestScopes]] @@ -517,9 +433,10 @@ def scope_cache(f: _GetScopesMethod) -> _GetScopesMethod: implementations that persists the resolved scopes on ``request.meta``. The readers of the resolved scopes — the synchronous readiness API of a - :ref:`throttling-aware scheduler ` and the - backoff methods — read this persisted value instead of resolving the scopes - again, so they stay cheap and consistent, and it survives a request being + :ref:`throttling-aware scheduler ` and + :meth:`~ThrottlingManagerProtocol.get_resolved_scopes` — read this persisted + value instead of resolving the scopes again, so they stay cheap and + consistent, and it survives a request being serialized to and restored from a :ref:`disk queue ` (which is what lets the readiness API resolve the scopes of a restored request synchronously). @@ -570,12 +487,6 @@ class ThrottlingManager: def __init__(self, crawler: Crawler) -> None: self.crawler = crawler - 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") # When set, each request also gets an IP scope (its resolved address), # enforced alongside its domain scope (see _resolve_scopes_sync). @@ -599,27 +510,6 @@ class ThrottlingManager: self._scopes_config: dict[str, dict[str, Any]] = crawler.settings.getdict( "THROTTLING_SCOPES" ) - # Cheap pre-filter for the backoff methods: the union of the global - # triggers and every per-scope override. A response whose status (or an - # exception type) is not in the union cannot trigger backoff for any - # scope, so its scopes need not be resolved (see get_response_backoff - # and get_exception_backoff). Each scope still makes the final decision - # via its own triggers_backoff_* methods. - self._any_backoff_http_codes: set[int] = set(self._backoff_http_codes) - self._any_backoff_exceptions: tuple[type[BaseException], ...] = ( - self._backoff_exceptions - ) - for scope_config in self._scopes_config.values(): - backoff = scope_config.get("backoff") or {} - if "http_codes" in backoff: - self._any_backoff_http_codes.update( - int(code) for code in backoff["http_codes"] - ) - if "exceptions" in backoff: - self._any_backoff_exceptions += tuple( - load_object(exc) if isinstance(exc, str) else exc - for exc in backoff["exceptions"] - ) # Ordered by least-recently-used first (see get_scope_manager), so the # scope limit can evict the coldest idle scopes (see THROTTLING_SCOPE_LIMIT). self._scope_managers: OrderedDict[ScopeID, ThrottlingScopeManagerProtocol] = ( @@ -675,88 +565,28 @@ class ThrottlingManager: scope_ids = sorted(iter_scopes(scopes)) return "+".join(scope_ids) if scope_ids else "" - def _cached_scope_values( - self, request: Request - ) -> list[tuple[ScopeID, float | None]]: - """Return the ``(scope_id, quota_amount)`` pairs of *request*, reading - the scopes persisted on ``request.meta`` by :meth:`get_scopes` (via - :func:`scope_cache`) and falling back to :meth:`_resolve_scopes_sync`. + def get_resolved_scopes(self, request: Request) -> RequestScopes: + """Return the scopes under which *request* was (or will be) sent, + reusing those persisted on ``request.meta`` by an earlier + :meth:`get_scopes` call (see :func:`scope_cache`) and falling back to + :meth:`_resolve_scopes_sync` only when none were persisted. - The persisted value is authoritative here: a request reaching the - readiness API has already been enqueued and filed under the queue for - these very scopes. - """ - if _RESOLVED_SCOPES_META_KEY in request.meta: - scopes = cast("RequestScopes", request.meta[_RESOLVED_SCOPES_META_KEY]) - else: - scopes = self._resolve_scopes_sync(request) - return list(iter_scope_values(scopes)) - - async def _scopes(self, request: Request) -> RequestScopes: - """Return the scopes of *request*, reusing those persisted on - ``request.meta`` by an earlier :meth:`get_scopes` call (see - :func:`scope_cache`) and resolving them only as a fallback. - - Used by the backoff methods so that a request whose scopes were already - resolved (at enqueue time, or by :meth:`acquire`) is not resolved again - — which, for a request restored from a disk queue, would also risk - attributing the backoff to different scopes than the ones it was sent - under. + The persisted value is authoritative: a request reaching a reader of + this method (the readiness API, or a component reacting to its response) + has already been enqueued and, for a request restored from a disk queue, + re-resolving could attribute it to different scopes than the ones it was + sent under. """ if _RESOLVED_SCOPES_META_KEY in request.meta: return cast("RequestScopes", request.meta[_RESOLVED_SCOPES_META_KEY]) - return await self.get_scopes(request) + return self._resolve_scopes_sync(request) - async def get_initial_backoff(self) -> BackoffData: - return None - - async def get_response_backoff(self, response: Response) -> BackoffData: - assert response.request is not None - if response.request.meta.get("throttling_dont_track"): - return None - if response.status not in self._any_backoff_http_codes: - return None - scopes = await self._scopes(response.request) - matched = [ - scope - for scope in iter_scopes(scopes) - if self.get_scope_manager(scope).triggers_backoff_for_status( - response.status - ) - ] - if not matched: - return None - if delay := self.get_response_delay(response): - return {scope: {"delay": delay} for scope in matched} - return matched - - def get_response_delay(self, response: Response) -> float | None: - """Return the throttling delay requested by the response.""" - retry_after = _parse_retry_after(response) - ratelimit_reset = _parse_ratelimit_reset(response) - if retry_after is None and ratelimit_reset is None: - return None - if retry_after is not None and ratelimit_reset is not None: - return max(retry_after, ratelimit_reset) - if retry_after is not None: - return retry_after - assert ratelimit_reset is not None - return ratelimit_reset - - async def get_exception_backoff( - self, request: Request, exception: Exception - ) -> BackoffData: - if request.meta.get("throttling_dont_track"): - return None - if not isinstance(exception, self._any_backoff_exceptions): - return None - scopes = await self._scopes(request) - matched = [ - scope - for scope in iter_scopes(scopes) - if self.get_scope_manager(scope).triggers_backoff_for_exception(exception) - ] - return matched or None + def _cached_scope_values( + self, request: Request + ) -> list[tuple[ScopeID, float | None]]: + """Return the ``(scope_id, quota_amount)`` pairs of *request*, from the + scopes returned by :meth:`get_resolved_scopes`.""" + return list(iter_scope_values(self.get_resolved_scopes(request))) # -- Scope-state coordination (called from the request lifecycle) -------- @@ -948,42 +778,29 @@ class ThrottlingManager: logger.debug(f"Holding {request} for {delay:.2f}s (throttling_delay)") return deadline - 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 + def back_off( + self, + scopes: RequestScopes, + *, + delay: float | None = None, + cap: bool = True, + ) -> None: + for scope_id in iter_scopes(scopes): if self._debug: logger.debug(f"Backoff for scope {scope_id} (delay: {delay})") - manager.record_backoff(delay=delay) + self.get_scope_manager(scope_id).record_backoff(delay=delay, cap=cap) + + def reconcile_quota( + self, + scopes: RequestScopes, + *, + consumed: float | None = None, + remaining: float | None = None, + ) -> None: + for scope_id in iter_scopes(scopes): + self.get_scope_manager(scope_id).reconcile_quota( + consumed=consumed, remaining=remaining + ) def _on_robots_parsed(self, robotparser: Any, request: Request) -> None: """Honor a robots.txt ``Crawl-delay`` on the :signal:`robots_parsed` @@ -1049,11 +866,9 @@ class ThrottlingManager: manager.set_concurrency(1) def delay_scope(self, scope_id: ScopeID, delay: float) -> None: - if self._debug: - logger.debug(f"Delaying scope {scope_id} for {delay:.2f}s") # Like a Retry-After / RateLimit-Reset header, this is a hard minimum # delay; unlike those, it is trusted, so it bypasses BACKOFF_MAX_DELAY. - self.get_scope_manager(scope_id).record_backoff(delay=float(delay), cap=False) + self.back_off(scope_id, delay=float(delay), cap=False) def get_scope_delay(self, scope_id: ScopeID) -> float: return self.get_scope_manager(scope_id).get_base_delay()