diff --git a/scrapy/throttling.py b/scrapy/throttling.py index e603f1fbe..bdddd3d3b 100644 --- a/scrapy/throttling.py +++ b/scrapy/throttling.py @@ -396,15 +396,28 @@ _GetScopesMethod = TypeVar( ) -def scope_cache(f: _GetScopesMethod) -> _GetScopesMethod: - """Decorator to cache the result of - :meth:`~ThrottlingManagerProtocol.get_scopes` calls. +# Request.meta key under which scope_cache persists the resolved scopes so that +# they survive a request being serialized to and restored from a disk queue. +_RESOLVED_SCOPES_META_KEY = "_throttling_resolved_scopes" - It should be used so that calls to - :meth:`~ThrottlingManagerProtocol.get_scopes` from methods like - :meth:`~ThrottlingManagerProtocol.get_response_backoff` or - :meth:`~ThrottlingManagerProtocol.get_exception_backoff` do not become - unnecessarily expensive. + +def scope_cache(f: _GetScopesMethod) -> _GetScopesMethod: + """Decorator for :meth:`~ThrottlingManagerProtocol.get_scopes` + 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 + serialized to and restored from a :ref:`disk queue ` (which is + what lets the readiness API resolve the scopes of a restored request + synchronously). + + The decorated method always re-resolves; it never reads the persisted value + back. So a request that inherited ``request.meta`` from another one (e.g. a + redirect built with :meth:`Request.replace() `, + which copies ``meta``) resolves its own scopes and overwrites the inherited + ones rather than reusing them. For example: @@ -422,15 +435,12 @@ def scope_cache(f: _GetScopesMethod) -> _GetScopesMethod: @wraps(f) async def wrapper(self: Any, request: Request) -> RequestScopes: - cache: WeakKeyDictionary[Request, RequestScopes] | None = self.__dict__.get( - "_scope_cache" - ) - if cache is None: - cache = self.__dict__["_scope_cache"] = WeakKeyDictionary() - if request in cache: - return cache[request] scopes = await f(self, request) - cache[request] = scopes + # Materialize one-shot iterables so the persisted value stays + # re-iterable and serializable. + if not isinstance(scopes, (str, dict)) and isinstance(scopes, Iterable): + scopes = list(scopes) + request.meta[_RESOLVED_SCOPES_META_KEY] = scopes return scopes return wrapper # type: ignore[return-value] @@ -476,7 +486,6 @@ class ThrottlingManager: "THROTTLING_SCOPES" ) self._scope_managers: dict[ScopeID, ThrottlingScopeManagerProtocol] = {} - self._global_concurrency: int = crawler.settings.getint("CONCURRENT_REQUESTS") self._last_eviction: float | None = None # Concurrency slots reserved by acquire(), to be released once the # request finishes downloading. @@ -491,12 +500,13 @@ class ThrottlingManager: def _resolve_scopes_sync(self, request: Request) -> RequestScopes: """Best-effort synchronous scope resolution. - It mirrors :meth:`get_scopes`, and is used by the synchronous + It backs :meth:`get_scopes` and is also the fallback for the synchronous readiness methods (:meth:`is_ready`, :meth:`reserve`, - :meth:`time_until_ready`) when the cached result of an earlier - :meth:`get_scopes` call is not available (e.g. for a request restored - from disk). Subclasses whose :meth:`get_scopes` cannot be resolved - synchronously should rely on the enqueue-time cache instead. + :meth:`time_until_ready`) when no scopes were persisted on + ``request.meta`` by an earlier :meth:`get_scopes` call (which normally + happens at enqueue time and survives disk restores; see + :func:`scope_cache`). Subclasses whose :meth:`get_scopes` cannot be + resolved synchronously rely on that persisted value instead. """ scopes = request.meta.get("throttling_scopes") if scopes is not None: @@ -507,16 +517,34 @@ class ThrottlingManager: self, request: Request ) -> list[tuple[ScopeID, float | None]]: """Return the ``(scope_id, quota_amount)`` pairs of *request*, reading - the cache populated by :meth:`get_scopes` and falling back to - :meth:`_resolve_scopes_sync`.""" - cache: WeakKeyDictionary[Request, RequestScopes] | None = self.__dict__.get( - "_scope_cache" - ) - scopes = cache.get(request) if cache is not None else None - if scopes is None: + the scopes persisted on ``request.meta`` by :meth:`get_scopes` (via + :func:`scope_cache`) and falling back to :meth:`_resolve_scopes_sync`. + + 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. + """ + if _RESOLVED_SCOPES_META_KEY in request.meta: + return cast("RequestScopes", request.meta[_RESOLVED_SCOPES_META_KEY]) + return await self.get_scopes(request) + async def get_initial_backoff(self) -> BackoffData: return None @@ -526,7 +554,7 @@ class ThrottlingManager: return None if response.status not in self._backoff_http_codes: return None - scopes = await self.get_scopes(response.request) + scopes = await self._scopes(response.request) if delay := self.get_response_delay(response): scopes = {scope: {"delay": delay} for scope in iter_scopes(scopes)} return scopes @@ -550,7 +578,7 @@ class ThrottlingManager: if request.meta.get("throttling_dont_track"): return None if isinstance(exception, self._backoff_exceptions): - return await self.get_scopes(request) + return await self._scopes(request) return None # -- Scope-state coordination (called from the request lifecycle) -------- @@ -635,6 +663,11 @@ class ThrottlingManager: return True def reserve(self, request: Request) -> None: + # A throttling-aware scheduler reserves every request before handing it + # to the engine, so acquire() always returns early for it and never + # gets to evict idle scopes; do it here so their managers do not pile + # up on broad crawls. + self._maybe_evict(time.monotonic()) managers = [ (self._get_scope_manager(scope_id), value) for scope_id, value in self._cached_scope_values(request) @@ -652,13 +685,7 @@ class ThrottlingManager: return wait if wait > 0 else None def scope_load(self, scope_id: ScopeID) -> float: - manager = self._get_scope_manager(scope_id) - active: int = getattr(manager, "_active", 0) - concurrency: int | None = getattr(manager, "_concurrency", None) - limit = concurrency if concurrency is not None else self._global_concurrency - if not limit: - return 0.0 - return active / limit + return self._get_scope_manager(scope_id).load() def get_request_delay(self, request: Request, now: float | None = None) -> float: now = time.monotonic() if now is None else now @@ -933,6 +960,19 @@ class ThrottlingScopeManagerProtocol(Protocol): Return ``False`` when no concurrency limit is enforced. """ + def load(self) -> float: + """Return the current load of this scope: a non-negative number, with + ``1.0`` meaning "as busy as its concurrency limit allows". + + A :ref:`throttling-aware scheduler ` uses + this to break ties between equally-prioritized requests, preferring the + least-loaded scopes. The reference implementation returns active sends + divided by the concurrency limit (falling back to + :setting:`CONCURRENT_REQUESTS` when the scope enforces no explicit + limit), but any consistent busyness metric works; return ``0.0`` when + none is meaningful. + """ + def slot_event(self) -> Deferred[None]: """Return a :class:`~twisted.internet.defer.Deferred` that fires when a concurrency slot next becomes available (e.g. when @@ -1042,6 +1082,9 @@ class ThrottlingScopeManager: self._concurrency = self._min_concurrency else: self._concurrency = None + # Used as the load denominator when the scope enforces no explicit + # concurrency limit (see load()). + self._global_concurrency: int = settings.getint("CONCURRENT_REQUESTS") # Quota. quota = config.get("quota") @@ -1090,6 +1133,14 @@ class ThrottlingScopeManager: def _recover(self, now: float) -> None: if self._backoff_level == 0 or self._last_backoff_time is None: return + if self._window <= 0: + # A non-positive window has no recovery cadence to step through (and + # would spin forever on a zero-length step), so recover at once. + self._backoff_level = 0 + self._delay = self._base_delay + self._in_backoff_until = None + self._last_backoff_time = 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 @@ -1105,6 +1156,10 @@ class ThrottlingScopeManager: under :setting:`RAMPUP_BACKOFF_TARGET` backoff triggers.""" if not self._rampup_enabled: return + if self._window <= 0: + # No window means no cadence to ramp up on (and a zero-length step + # would spin forever). + return if self._rampup_window_start is None: self._rampup_window_start = now return @@ -1131,6 +1186,12 @@ class ThrottlingScopeManager: def _maybe_reset_quota(self, now: float) -> None: if self._quota is None: return + if self._quota_window <= 0: + # A non-positive window has no reset cadence to step through (and + # would spin forever on a zero-length step), so keep it reset. + self._consumed = 0.0 + self._quota_window_start = now + return if self._quota_window_start is None: self._quota_window_start = now return @@ -1180,6 +1241,16 @@ class ThrottlingScopeManager: def concurrency_blocked(self) -> bool: return self._concurrency is not None and self._active >= self._concurrency + def load(self) -> float: + limit = ( + self._concurrency + if self._concurrency is not None + else self._global_concurrency + ) + if not limit: + return 0.0 + return self._active / limit + 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 diff --git a/tests/test_throttling.py b/tests/test_throttling.py index 3dc7727f9..9ba6d51eb 100644 --- a/tests/test_throttling.py +++ b/tests/test_throttling.py @@ -17,9 +17,11 @@ from scrapy.throttling import ( add_scope, iter_scope_values, iter_scopes, + scope_cache, update_scope_backoff, ) from scrapy.utils.defer import deferred_from_coro, maybe_deferred_to_future +from scrapy.utils.httpobj import urlparse_cached from scrapy.utils.test import get_crawler from tests.spiders import SimpleSpider from tests.utils.decorators import coroutine_test @@ -69,8 +71,8 @@ class TestThrottlingManager: 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). + # get_scopes is deterministic per request, so a second call yields the + # same scopes. assert await manager.get_scopes(request) == first @coroutine_test @@ -87,6 +89,80 @@ class TestThrottlingManager: ) assert await manager.get_scopes(request) == {"api": 2.0} + @coroutine_test + async def test_get_scopes_persisted_in_meta(self): + from scrapy.throttling import _RESOLVED_SCOPES_META_KEY # noqa: PLC0415 + + manager = _manager() + request = Request("http://example.com/a") + scopes = await manager.get_scopes(request) + assert request.meta[_RESOLVED_SCOPES_META_KEY] == scopes + + @coroutine_test + async def test_scope_cache_works_without_crawler(self): + # scope_cache only persists to meta; it needs nothing from the manager. + from scrapy.throttling import _RESOLVED_SCOPES_META_KEY # noqa: PLC0415 + + class CrawlerlessManager: + @scope_cache + async def get_scopes(self, request): + return "scope" + + request = Request("http://example.com/a") + assert await CrawlerlessManager().get_scopes(request) == "scope" + assert request.meta[_RESOLVED_SCOPES_META_KEY] == "scope" + + @coroutine_test + async def test_backoff_reuses_persisted_scopes(self): + # Once get_scopes has resolved and persisted the scopes, the backoff + # path reuses them instead of resolving again. + calls = [] + + class CountingManager(ThrottlingManager): + @scope_cache + async def get_scopes(self, request): + calls.append(request.url) + return urlparse_cached(request).netloc + + manager = CountingManager( + get_crawler(settings_dict={"BACKOFF_EXCEPTIONS": ["builtins.ValueError"]}) + ) + request = Request("http://example.com/a") + await manager.get_scopes(request) + assert calls == ["http://example.com/a"] + assert ( + await manager.get_exception_backoff(request, ValueError()) == "example.com" + ) + # No second resolution. + assert calls == ["http://example.com/a"] + + @coroutine_test + async def test_get_scopes_survives_disk_roundtrip(self): + from scrapy.utils.request import request_from_dict # noqa: PLC0415 + + manager = _manager() + request = Request( + "http://example.com/a", meta={"throttling_scopes": {"bucket": 3.0}} + ) + await manager.get_scopes(request) + # A request restored from a disk queue is a fresh object; the synchronous + # readiness path must still recover its scopes (with quota values) from + # the persisted meta, without re-running get_scopes. + restored = request_from_dict(request.to_dict()) + assert manager._cached_scope_values(restored) == [("bucket", 3.0)] + + @coroutine_test + async def test_get_scopes_reresolved_after_cross_host_replace(self): + # A redirect built with Request.replace() copies meta (including the + # persisted scopes), but get_scopes always re-resolves (it never reads + # the persisted value back), so it must not reuse the original host's + # scopes. + manager = _manager() + request = Request("http://example.com/a") + assert await manager.get_scopes(request) == "example.com" + redirected = request.replace(url="http://other.example/a") + assert await manager.get_scopes(redirected) == "other.example" + @coroutine_test async def test_get_initial_backoff_none(self): manager = _manager() @@ -324,6 +400,21 @@ class TestThrottlingManager: # even though it has been idle for longer than THROTTLING_SCOPE_MAX_IDLE. assert "example.com" in manager._scope_managers + def test_reserve_evicts_idle_scopes(self): + # A throttling-aware scheduler reserves every request before the engine + # reaches acquire() (which fast-paths reserved requests), so reserve() + # must be the hook that evicts idle scope managers; otherwise they pile + # up unbounded on broad crawls. + manager = _manager({"THROTTLING_SCOPE_MAX_IDLE": 1.0}) + idle = manager._get_scope_manager("idle.example") + # Make it look long-idle: a finished send in the distant monotonic past. + idle.record_sent(now=0.0) + idle.record_done(now=0.0) + assert "idle.example" in manager._scope_managers + manager.reserve(Request("http://active.example/1")) + assert "idle.example" not in manager._scope_managers + assert "active.example" in manager._scope_managers + class TestThrottlingScopeManager: def test_no_delay_by_default(self): @@ -467,6 +558,34 @@ class TestThrottlingScopeManager: scope.record_done(now=0.0) assert event.called + def test_zero_backoff_window_recovers_at_once(self): + # A non-positive window must not make _recover spin forever; it recovers + # fully on the next can_send() instead. + scope = _scope_manager(settings={"BACKOFF_WINDOW": 0}, config={"id": "x"}) + scope.record_backoff(now=0.0) + assert scope._backoff_level == 1 + scope.can_send(now=10.0) + assert scope._backoff_level == 0 + assert scope._delay == scope._base_delay + + def test_zero_quota_window_keeps_quota_reset(self): + # A non-positive quota window must not make _maybe_reset_quota spin + # forever; the quota stays continuously reset instead. + scope = _scope_manager(config={"id": "x", "quota": 10.0, "window": 0}) + scope.record_sent(now=0.0, amount=10.0) + assert scope.can_send(now=1.0, amount=5.0) == 0.0 + assert scope._consumed == 0.0 + + def test_zero_window_disables_rampup(self): + # A non-positive window must not make _maybe_rampup spin forever; rampup + # is simply disabled. + scope = _scope_manager( + settings={"BACKOFF_WINDOW": 0}, config={"id": "x", "rampup": True} + ) + scope.can_send(now=0.0) + scope.can_send(now=10_000.0) + assert scope._concurrency == scope._min_concurrency + def test_set_concurrency_fires_slot_event(self): scope = _scope_manager(config={"id": "x", "concurrency": 1}) scope.record_sent(now=0.0)