From 4ac921f3891268b96bed711ddc75eff0e56c958f Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Adri=C3=A1n=20Chaves?= Date: Mon, 24 Mar 2025 17:26:58 +0100 Subject: [PATCH] WIP --- docs/news.rst | 18 -- docs/topics/request-response.rst | 9 +- docs/topics/settings.rst | 23 -- docs/topics/spider-middleware.rst | 17 -- docs/topics/spiders.rst | 51 ++++ scrapy/__init__.py | 2 - scrapy/commands/shell.py | 2 - scrapy/core/_seeding.py | 68 ------ scrapy/core/engine.py | 93 ++------ scrapy/http/request/__init__.py | 12 +- scrapy/settings/default_settings.py | 4 - scrapy/spiders/__init__.py | 16 +- tests/test_crawl.py | 20 +- tests/test_engine_loop.py | 130 ++++++++++ tests/test_engine_seeding.py | 358 ---------------------------- tests/test_scheduler.py | 35 ++- 16 files changed, 253 insertions(+), 605 deletions(-) delete mode 100644 scrapy/core/_seeding.py create mode 100644 tests/test_engine_loop.py delete mode 100644 tests/test_engine_seeding.py diff --git a/docs/news.rst b/docs/news.rst index b88626bd5..37ed1e693 100644 --- a/docs/news.rst +++ b/docs/news.rst @@ -19,9 +19,6 @@ Backward-incompatible changes - By default, the iteration of start requests and items no longer stops once there are requests in the scheduler. - You can restore the previous behavior by setting :setting:`SEEDING_POLICY` - to :py:enum:mem:`~scrapy.SeedingPolicy.lazy`. - - In ``scrapy.core.engine.ExecutionEngine``: - The second parameter of ``open_spider()``, ``start_requests()``, has @@ -69,21 +66,6 @@ New features (:issue:`456`, :issue:`3477`, :issue:`4467`, :issue:`5627`, :issue:`6729`) -- The new :setting:`SEEDING_POLICY` setting allows customizing how start - requests and items are iterated. - - You can also override the active seeding policy from - :meth:`Spider.start ` and from - :meth:`SpiderMiddleware.process_start - `. - - .. note:: Some third-party spider middlewares may need to be updated for - Scrapy VERSION support before you can use them in combination with the - ability to override the active seeding policy. - - (:issue:`740`, :issue:`1051`, :issue:`1443`, :issue:`3237`, :issue:`4467`, - :issue:`5282`, :issue:`6729`) - Bug fixes ~~~~~~~~~ diff --git a/docs/topics/request-response.rst b/docs/topics/request-response.rst index 2032fc5cf..de7840878 100644 --- a/docs/topics/request-response.rst +++ b/docs/topics/request-response.rst @@ -127,10 +127,7 @@ Request objects body to bytes (if given as a string). :type encoding: str - :param priority: the priority of this request (defaults to ``0``). - The priority is used by the scheduler to define the order used to process - requests. Requests with a higher priority value will execute earlier. - Negative values are allowed in order to indicate relatively low-priority. + :param priority: sets :attr:`priority`, defaults to ``0``. :type priority: int :param dont_filter: sets :attr:`dont_filter`, defaults to ``False``. @@ -179,6 +176,8 @@ Request objects .. autoattribute:: errback + .. autoattribute:: priority + .. attribute:: Request.cb_kwargs A dictionary that contains arbitrary metadata for this request. Its contents @@ -353,7 +352,7 @@ errors if needed: "https://example.invalid/", # DNS error expected ] - async def start(self): + async def yield_seeds(self): for u in self.start_urls: yield scrapy.Request( u, diff --git a/docs/topics/settings.rst b/docs/topics/settings.rst index 261100e00..67de4f5f1 100644 --- a/docs/topics/settings.rst +++ b/docs/topics/settings.rst @@ -1733,29 +1733,6 @@ Soft limit (in bytes) for response data being processed. While the sum of the sizes of all responses being processed is above this value, Scrapy does not process new requests. -.. setting:: SEEDING_POLICY - -SEEDING_POLICY --------------- - -.. versionadded:: VERSION - -Default: :py:enum:mem:`SeedingPolicy.greedy ` - -Determines the way :meth:`Spider.start ` is -iterated. - -Its value may be defined as a member of the :class:`~scrapy.SeedingPolicy` enum -(e.g. :py:enum:mem:`SeedingPolicy.lazy `) or as a -matching string (e.g. ``"lazy"``). - -You can also override the active seeding policy from :meth:`Spider.start -` and from :meth:`SpiderMiddleware.process_start -`. - -.. autoenum:: scrapy.SeedingPolicy - :members: - .. setting:: SPIDER_CONTRACTS SPIDER_CONTRACTS diff --git a/docs/topics/spider-middleware.rst b/docs/topics/spider-middleware.rst index 79b8b1555..0baf094bc 100644 --- a/docs/topics/spider-middleware.rst +++ b/docs/topics/spider-middleware.rst @@ -85,23 +85,6 @@ one or more of these methods: You may yield the same type of objects as :meth:`~scrapy.Spider.start`. - As with :meth:`~scrapy.Spider.start`, how this method is iterated by - default is controlled by :setting:`SEEDING_POLICY`. It is also possible - to yield a :class:`~scrapy.SeedingPolicy` enum or a matching string to - change the active seeding policy, for example: - - .. code-block:: python - - async def process_start(self, start): - yield "front_load" - async for item_or_request in start: - yield item_or_request - yield "idle" - - .. tip:: You can also restore the configured seeding policy by - :ref:`reading its value ` from the - :setting:`SEEDING_POLICY` setting and yielding it. - To write spider middlewares that work on Scrapy versions lower than VERSION, define also a synchronous ``process_start_requests()`` method that returns an iterable. For example: diff --git a/docs/topics/spiders.rst b/docs/topics/spiders.rst index 71775247c..2432223eb 100644 --- a/docs/topics/spiders.rst +++ b/docs/topics/spiders.rst @@ -364,6 +364,57 @@ used by :class:`~scrapy.downloadermiddlewares.useragent.UserAgentMiddleware`:: Spider arguments can also be passed through the Scrapyd ``schedule.json`` API. See `Scrapyd documentation`_. +.. _spider-start: + +Spider start +============ + +The way :meth:`~scrapy.Spider.start` works may be counterintuitive. + +By default, if you do not use :ref:`await ` in +:meth:`~scrapy.Spider.start`, or if you instead use the +:attr:`~scrapy.Spider.start_urls` attribute, this is what happens: + +- The first 8-16 start requests (based on :setting:`CONCURRENT_REQUESTS` and + :setting:`CONCURRENT_REQUESTS_PER_DOMAIN`) are sent in the order in which + they are yielded. + + .. note:: Responses may come in a different order. + +- The remaining start requests are sent in reverse order, and only when there + are not enough pending requests yielded from callbacks to reach the + configured concurrency. + + This is because the :ref:`scheduler `, where pending + requests are stored, uses a LIFO (last in, first out) queue by default, + configured in the :setting:`SCHEDULER_MEMORY_QUEUE` and + :setting:`SCHEDULER_DISK_QUEUE`. + + So, provided all pending requests have the same + :attr:`~scrapy.Request.priority`, scheduled requests are sent in reserve + order. The first few requests are sent in order only because they are sent + as soon as they are scheduled. + +If you need start requests to be sent before requests yielded from spider +callbacks, you can set a higher priority for them. For example: + +.. code-block:: python + + async def start(self): + async for request in super().start(): + yield request.replace(priority=1) + +If you also need them to be sent in order, you can assign them decreasing +priority values. For example: + +.. code-block:: python + + async def start(self): + priority = len(self.start_urls) + async for request in super().start(): + yield request.replace(priority=priority) + priority -= 1 + .. _builtin-spiders: Generic Spiders diff --git a/scrapy/__init__.py b/scrapy/__init__.py index 7bb958a2d..256504c9c 100644 --- a/scrapy/__init__.py +++ b/scrapy/__init__.py @@ -7,7 +7,6 @@ import sys import warnings # Declare top-level shortcuts -from scrapy.core._seeding import SeedingPolicy from scrapy.http import FormRequest, Request from scrapy.item import Field, Item from scrapy.selector import Selector @@ -18,7 +17,6 @@ __all__ = [ "FormRequest", "Item", "Request", - "SeedingPolicy", "Selector", "Spider", "__version__", diff --git a/scrapy/commands/shell.py b/scrapy/commands/shell.py index c50c963ef..4c3aca605 100644 --- a/scrapy/commands/shell.py +++ b/scrapy/commands/shell.py @@ -9,7 +9,6 @@ from __future__ import annotations from threading import Thread from typing import TYPE_CHECKING, Any -from scrapy import SeedingPolicy from scrapy.commands import ScrapyCommand from scrapy.http import Request from scrapy.shell import Shell @@ -28,7 +27,6 @@ class Command(ScrapyCommand): "DUPEFILTER_CLASS": "scrapy.dupefilters.BaseDupeFilter", "KEEP_ALIVE": True, "LOGSTATS_INTERVAL": 0, - "SEEDING_POLICY": SeedingPolicy.lazy, } def syntax(self) -> str: diff --git a/scrapy/core/_seeding.py b/scrapy/core/_seeding.py deleted file mode 100644 index 08c68951c..000000000 --- a/scrapy/core/_seeding.py +++ /dev/null @@ -1,68 +0,0 @@ -from enum import Enum - -try: - from enum_tools.documentation import document_enum -except ImportError: - - def document_enum(func): # type: ignore[misc] - return func -else: - # https://github.com/domdfcoding/enum_tools/issues/29 - import enum_tools.documentation - - enum_tools.documentation.INTERACTIVE = True - - -@document_enum -class SeedingPolicy(Enum): - front_load = "front_load" - """The crawl does not start until all start requests have been scheduled. - - Aims to give the :ref:`scheduler ` full control over - request order from the start. Some custom schedulers may require this - seeding policy to work as designed. - """ - - greedy = "greedy" - """Iterating start items and requests takes priority over processing - scheduled requests. - - Every time a start request is iterated, it is scheduled, and then the next - request from the scheduler is sent. - - .. note:: That request sent may not be the scheduled start request - depending on the priority of scheduled requests, on the configured - :setting:`SCHEDULER` and on certain scheduler settings (e.g. - :setting:`SCHEDULER_MEMORY_QUEUE`). - - Best used when prioritizing start requests is important. - """ - - idle = "idle" - """A single start item or request is read only when there are neither - scheduled nor on-going requests. - - That is, a new start item or request is not read until all requests - triggered by the previous start request, directly or indirectly, have been - processed. - - Unlike :py:enum:mem:`lazy`, resource savings are prioritized over crawl - speed. - - It is functionally equivalent to running a spider multiple times in a row, - one per start request. - """ - - lazy = "lazy" - """Processing scheduled requests takes priority over iterating start items - and requests. - - Aims to minimize the number of requests in the scheduler at any given time, - to minimize resource usage (memory or disk, depending on - :setting:`JOBDIR`). - - It is best used when start request priority is not important. - - Switching to :py:enum:mem:`idle` may lower resource usage further at the - cost of also lowering crawl speed. - """ diff --git a/scrapy/core/engine.py b/scrapy/core/engine.py index b87fd1b78..6ec9780e6 100644 --- a/scrapy/core/engine.py +++ b/scrapy/core/engine.py @@ -24,8 +24,6 @@ from scrapy.utils.log import failure_to_exc_info, logformatter_adapter from scrapy.utils.misc import build_from_crawler, load_object from scrapy.utils.reactor import CallLaterOnce -from ._seeding import SeedingPolicy - if TYPE_CHECKING: from collections.abc import AsyncIterable, Callable, Generator @@ -43,10 +41,6 @@ logger = logging.getLogger(__name__) _T = TypeVar("_T") -class _SeedingPolicyChange(Exception): - pass - - class _Slot: def __init__( self, @@ -109,20 +103,7 @@ class ExecutionEngine: spider_closed_callback ) self.start_time: float | None = None - self._load_seeding_policy() self._start: AsyncIterable[Any] | None = None - self._waiting_for_seed: bool = False - - def _load_seeding_policy(self) -> None: - try: - self._seeding_policy = SeedingPolicy(self.settings["SEEDING_POLICY"]) - except ValueError: - supported_values = ", ".join(policy.value for policy in SeedingPolicy) - raise ValueError( - f"The value of the SEEDING_POLICY setting " - f"({self.settings['SEEDING_POLICY']!r}) is not supported. " - f"Supported values: {supported_values}." - ) def _get_scheduler_class(self, settings: BaseSettings) -> type[BaseScheduler]: from scrapy.core.scheduler import BaseScheduler @@ -185,10 +166,7 @@ class ExecutionEngine: self.paused = False @inlineCallbacks - def _process_next_seed(self): - if self._waiting_for_seed: - return - self._waiting_for_seed = True + def _process_next_spider_start_yield(self): try: item_or_request = yield deferred_from_coro(self._start.__anext__()) except StopAsyncIteration: @@ -203,68 +181,28 @@ class ExecutionEngine: else: if isinstance(item_or_request, Request): self.crawl(item_or_request) - if ( - self._seeding_policy is not SeedingPolicy.front_load - and not self._needs_backout() - ): - self._start_scheduled_request() - elif isinstance(item_or_request, (str, SeedingPolicy)): - try: - self._seeding_policy = SeedingPolicy(item_or_request) - except ValueError: - valid_policy_strings = ", ".join( - policy.value for policy in SeedingPolicy - ) - logger.error( - f"Start value {item_or_request!r} has been ignored. " - f"Start values of {str} type must be valid seeding " - f"policies ({valid_policy_strings})." - ) - self._slot.nextcall.schedule() - else: - raise _SeedingPolicyChange else: self.scraper.start_itemproc(item_or_request, response=None) self._slot.nextcall.schedule() - finally: - self._waiting_for_seed = False - if self._seeding_policy is SeedingPolicy.front_load and self._start is None: - self._slot.nextcall.schedule() @inlineCallbacks - def _start_next_requests(self) -> Generator[Deferred[Any], Any, None]: + def _process_spider_start(self) -> Generator[Deferred[Any], Any, None]: + """Process start items and requests in an asynchronous loop. + + Items are scraped. Requests are scheduled. + """ + while self._start is not None: + yield self._process_next_spider_start_yield() + if not self._needs_backout(): + self._slot.nextcall.schedule() + + def _start_next_requests(self) -> None: if self._slot is None or self._slot.closing is not None or self.paused: return - try: - if self._seeding_policy in {SeedingPolicy.idle, SeedingPolicy.lazy}: - while not self._needs_backout(): - if self._start_scheduled_request() is None: - break - if ( - self._start is not None - and not self._needs_backout() - and ( - self._seeding_policy is not SeedingPolicy.idle - or (not self._waiting_for_seed and not self.downloader.active) - ) - ): - yield self._process_next_seed() - else: - assert self._seeding_policy in { - SeedingPolicy.front_load, - SeedingPolicy.greedy, - } - if self._start is not None: - if not self._needs_backout(): - yield self._process_next_seed() - else: - while not self._needs_backout(): - if self._start_scheduled_request() is None: - break - except _SeedingPolicyChange: - self._slot.nextcall.schedule() - return + while not self._needs_backout(): + if self._start_scheduled_request() is None: + break if self.spider_is_idle() and self._slot.close_if_idle: self._spider_idle() @@ -452,6 +390,7 @@ class ExecutionEngine: assert self.crawler.stats self.crawler.stats.open_spider(spider) yield self.signals.send_catch_log_deferred(signals.spider_opened, spider=spider) + self._process_spider_start() self._slot.nextcall.schedule() self._slot.heartbeat.start(self._SLOT_HEARTBEAT_INTERVAL) diff --git a/scrapy/http/request/__init__.py b/scrapy/http/request/__init__.py index d72c72ecd..0bfc92c6b 100644 --- a/scrapy/http/request/__init__.py +++ b/scrapy/http/request/__init__.py @@ -130,6 +130,16 @@ class Request(object_ref): self._set_body(body) if not isinstance(priority, int): raise TypeError(f"Request priority not an integer: {priority!r}") + + #: Default: ``0`` + #: + #: Value that the :ref:`scheduler ` may use for + #: request prioritization. + #: + #: Built-in schedulers prioritize requests with a higher priority + #: value. + #: + #: Negative values are allowed. self.priority: int = priority if not (callable(callback) or callback is None): @@ -191,7 +201,7 @@ class Request(object_ref): #: #: When defining the start URLs of a spider through #: :attr:`~scrapy.Spider.start_urls`, this attribute is enabled by - #: default. See :meth:`~scrapy.Spider.start`. + #: default. See :meth:`~scrapy.Spider.yield_seeds`. self.dont_filter: bool = dont_filter self._meta: dict[str, Any] | None = dict(meta) if meta else None diff --git a/scrapy/settings/default_settings.py b/scrapy/settings/default_settings.py index 886154bfa..645e50301 100644 --- a/scrapy/settings/default_settings.py +++ b/scrapy/settings/default_settings.py @@ -17,8 +17,6 @@ import sys from importlib import import_module from pathlib import Path -from scrapy import SeedingPolicy - ADDONS = {} AJAXCRAWL_ENABLED = False @@ -310,8 +308,6 @@ SCHEDULER_PRIORITY_QUEUE = "scrapy.pqueues.ScrapyPriorityQueue" SCRAPER_SLOT_MAX_ACTIVE_SIZE = 5000000 -SEEDING_POLICY = SeedingPolicy.greedy - SPIDER_LOADER_CLASS = "scrapy.spiderloader.SpiderLoader" SPIDER_LOADER_WARN_ONLY = False diff --git a/scrapy/spiders/__init__.py b/scrapy/spiders/__init__.py index 6eba1b188..30e197b8c 100644 --- a/scrapy/spiders/__init__.py +++ b/scrapy/spiders/__init__.py @@ -113,20 +113,6 @@ class Spider(object_ref): async def start(self): yield {"foo": "bar"} - Use :setting:`SEEDING_POLICY` to set how :meth:`start` is - iterated by default. It is also - possible to yield a :class:`~scrapy.SeedingPolicy` enum or a matching - string to change the active seeding policy, for example: - - .. code-block:: python - - async def start(self): - yield "front_load" - yield Request("https://a.example") - yield Request("https://b.example") - yield self.crawler.settings["SEEDING_POLICY"] - yield Request("https://c.example") - To write spiders that work on Scrapy versions lower than VERSION, define also a synchronous ``start_requests()`` method that returns an iterable. For example: @@ -135,6 +121,8 @@ class Spider(object_ref): def start_requests(self): yield Request("https://toscrape.com/") + + .. seealso:: :ref:`spider-start` """ for item_or_request in self.start_requests(): yield item_or_request diff --git a/tests/test_crawl.py b/tests/test_crawl.py index 9115edd27..7793505da 100644 --- a/tests/test_crawl.py +++ b/tests/test_crawl.py @@ -194,25 +194,15 @@ class TestCrawl(TestCase): @defer.inlineCallbacks def test_start_unsupported_output(self): - """Anything that is not a request, a seeding policy or a string (which - is assumed to be a seeding policy) is assumed to be an item, avoiding a - potentially expensive call to itemadapter.is_item, and letting instead - things fail when ItemAdapter is actually used on the corresponding - non-item object.""" + """Anything that is not a request is assumed to be an item, avoiding a + potentially expensive call to itemadapter.is_item(), and letting + instead things fail when ItemAdapter is actually used on the + corresponding non-item object.""" with LogCapture("scrapy", level=logging.ERROR) as log: crawler = get_crawler(StartGoodAndBadOutput) yield crawler.crawl(mockserver=self.mockserver) - assert len(log.records) == 1 - - @defer.inlineCallbacks - def test_start_laziness(self): - settings = {"CONCURRENT_REQUESTS": 1, "SEEDING_POLICY": "lazy"} - crawler = get_crawler(BrokenStartSpider, settings) - yield crawler.crawl(mockserver=self.mockserver) - assert crawler.spider.seedsseen.index(None) < crawler.spider.seedsseen.index( - 99 - ), crawler.spider.seedsseen + assert len(log.records) == 0 @defer.inlineCallbacks def test_start_dupes(self): diff --git a/tests/test_engine_loop.py b/tests/test_engine_loop.py new file mode 100644 index 000000000..62e941f49 --- /dev/null +++ b/tests/test_engine_loop.py @@ -0,0 +1,130 @@ +from collections import deque + +from twisted.internet.defer import Deferred +from twisted.trial.unittest import TestCase + +from scrapy import Request, Spider, signals +from scrapy.core.engine import ExecutionEngine +from scrapy.utils.defer import deferred_f_from_coro_f, maybe_deferred_to_future +from scrapy.utils.test import get_crawler + +from .mockserver import MockServer +from .test_scheduler import MemoryScheduler + + +def sleep(seconds: float = 0.001): + from twisted.internet import reactor + + deferred: Deferred[None] = Deferred() + reactor.callLater(seconds, deferred.callback, None) + return maybe_deferred_to_future(deferred) + + +class MainTestCase(TestCase): + @deferred_f_from_coro_f + async def test_sleep(self): + """Neither asynchronous sleeps on Spider.start() nor the equivalent on + the scheduler (returning no requests while also returning True from + the has_pending_requests() method) should cause the spider to miss the + processing of any later requests.""" + seconds = ExecutionEngine._SLOT_HEARTBEAT_INTERVAL + 0.01 + + class TestSpider(Spider): + name = "test" + + async def start(self): + from twisted.internet import reactor + + yield Request("data:,a") + + await sleep(seconds) + + self.crawler.engine._slot.scheduler.pause() + self.crawler.engine._slot.scheduler.enqueue_request(Request("data:,b")) + + # During this time, the reactor reports having requests but + # returns None. + await sleep(seconds) + + self.crawler.engine._slot.scheduler.unpause() + + # The scheduler request is processed. + await sleep(seconds) + + yield Request("data:,c") + + await sleep(seconds) + + self.crawler.engine._slot.scheduler.pause() + self.crawler.engine._slot.scheduler.enqueue_request(Request("data:,d")) + + # The last start request is processed during the time until the + # delayed call below, proving that the start iteration can + # finish before a scheduler “sleep” without causing the + # scheduler to finish. + reactor.callLater(seconds, self.crawler.engine._slot.scheduler.unpause) + + def parse(self, response): + pass + + actual_urls = [] + + def track_url(request, spider): + actual_urls.append(request.url) + + settings = {"SCHEDULER": MemoryScheduler} + crawler = get_crawler(TestSpider, settings_dict=settings) + crawler.signals.connect(track_url, signals.request_reached_downloader) + await maybe_deferred_to_future(crawler.crawl()) + assert crawler.stats.get_value("finish_reason") == "finished" + expected_urls = ["data:,a", "data:,b", "data:,c", "data:,d"] + assert actual_urls == expected_urls, f"{actual_urls=} != {expected_urls=}" + + +class MockServerTestCase(TestCase): + # See the comment on the matching line above. + timeout = ExecutionEngine._SLOT_HEARTBEAT_INTERVAL + + @classmethod + def setUpClass(cls): + cls.mockserver = MockServer() + cls.mockserver.__enter__() + + @classmethod + def tearDownClass(cls): + cls.mockserver.__exit__(None, None, None) + + @deferred_f_from_coro_f + async def test_default_behavior(self): + """Verify the behavior that the docs claims is the default when it + comes to start request send order.""" + seconds = 0.1 + + def _url(id, delay=seconds): + return self.mockserver.url(f"/delay?n={seconds}&{id}") + + class TestSpider(Spider): + name = "test" + start_urls = [_url("a", delay=0), _url("b"), _url("d"), _url("e")] + queue = deque([Request(_url("c"))]) + + def parse(self, response): + try: + request = self.queue.popleft() + except IndexError: + pass + else: + yield request + + actual_urls = [] + + def track_url(request, spider): + actual_urls.append(request.url) + + settings = {"CONCURRENT_REQUESTS": 2} + crawler = get_crawler(TestSpider, settings_dict=settings) + crawler.signals.connect(track_url, signals.request_reached_downloader) + await maybe_deferred_to_future(crawler.crawl()) + assert crawler.stats.get_value("finish_reason") == "finished" + expected_urls = [_url(letter) for letter in "abcd"] + assert actual_urls == expected_urls, f"{actual_urls=} != {expected_urls=}" diff --git a/tests/test_engine_seeding.py b/tests/test_engine_seeding.py deleted file mode 100644 index c52a2b8bd..000000000 --- a/tests/test_engine_seeding.py +++ /dev/null @@ -1,358 +0,0 @@ -from __future__ import annotations - -from collections import defaultdict, deque -from logging import ERROR - -from testfixtures import LogCapture -from twisted.trial.unittest import TestCase - -from scrapy import Request, SeedingPolicy, Spider, signals -from scrapy.core.engine import ExecutionEngine -from scrapy.core.scheduler import BaseScheduler -from scrapy.utils.defer import deferred_f_from_coro_f, maybe_deferred_to_future -from scrapy.utils.test import get_crawler - -from .mockserver import MockServer -from .test_spider_start import twisted_sleep - - -class MainTestCase(TestCase): - # If the test ends before the heartbeat, it may mean that the logic to - # re-schecule a new call of _start_next_requests under the right - # ciscumstances is not properly implemented, and the hearatbeat is working - # as a workaround for that issue. This is a performance issue and should - # be addressed. - # - # It could also happen that, on some CI runners, some tests (e.g. those - # below using a mock server) run too slow and proper handling overlaps with - # the heartbeat. If that is the case, it may be worth considering - # increasing the heartbeat time. It should be safe, since in most real live - # scenarios the heartbeat should never make a difference, and we may - # eventually remove the heartbeat altogether. - timeout = ExecutionEngine._SLOT_HEARTBEAT_INTERVAL - - @deferred_f_from_coro_f - async def test_greedy(self): - class TestScheduler(BaseScheduler): - def __init__(self, *args, **kwargs): - self.requests = deque((Request("data:,b"),)) - - def enqueue_request(self, request: Request) -> bool: - self.requests.append(request) - return True - - def has_pending_requests(self) -> bool: - return bool(self.requests) - - def next_request(self) -> Request | None: - try: - return self.requests.pop() - except IndexError: - return None - - class TestSpider(Spider): - name = "test" - start_urls = ["data:,a"] - - def parse(self, response): - pass - - actual_urls = [] - - def track_url(request, spider): - actual_urls.append(request.url) - - settings = {"SCHEDULER": TestScheduler} - crawler = get_crawler(TestSpider, settings_dict=settings) - crawler.signals.connect(track_url, signals.request_reached_downloader) - await maybe_deferred_to_future(crawler.crawl()) - assert crawler.stats.get_value("finish_reason") == "finished" - expected_urls = ["data:,a", "data:,b"] - assert actual_urls == expected_urls, f"{actual_urls=} != {expected_urls=}" - - @deferred_f_from_coro_f - async def test_lazy(self): - class TestScheduler(BaseScheduler): - def __init__(self, *args, **kwargs): - self.requests = deque((Request("data:,a"),)) - - def enqueue_request(self, request: Request) -> bool: - self.requests.append(request) - return True - - def has_pending_requests(self) -> bool: - return bool(self.requests) - - def next_request(self) -> Request | None: - try: - return self.requests.popleft() - except IndexError: - return None - - class TestSpider(Spider): - name = "test" - start_urls = ["data:,b"] - - def parse(self, response): - pass - - actual_urls = [] - - def track_url(request, spider): - actual_urls.append(request.url) - - settings = {"SCHEDULER": TestScheduler, "SEEDING_POLICY": "lazy"} - crawler = get_crawler(TestSpider, settings_dict=settings) - crawler.signals.connect(track_url, signals.request_reached_downloader) - await maybe_deferred_to_future(crawler.crawl()) - assert crawler.stats.get_value("finish_reason") == "finished" - expected_urls = ["data:,a", "data:,b"] - assert actual_urls == expected_urls, f"{actual_urls=} != {expected_urls=}" - - @deferred_f_from_coro_f - async def test_lazy_blocking(self): - """If the scheduler reports having requests but yields none, the lazy - policy schedules start requests.""" - - class TestScheduler(BaseScheduler): - def __init__(self, *args, **kwargs): - self.requests = deque() - self.stop = False - - def enqueue_request(self, request: Request) -> bool: - self.requests.append(request) - return True - - def has_pending_requests(self) -> bool: - return not self.stop - - def next_request(self) -> Request | None: - try: - return self.requests.popleft() - except IndexError: - return None - - sleep_seconds = 0.0001 - - class TestSpider(Spider): - name = "test" - - async def start(self): - await maybe_deferred_to_future(twisted_sleep(sleep_seconds)) - yield Request("data:,a") - await maybe_deferred_to_future(twisted_sleep(sleep_seconds)) - self.crawler.engine._slot.scheduler.enqueue_request(Request("data:,b")) - await maybe_deferred_to_future(twisted_sleep(sleep_seconds)) - yield Request("data:,c") - self.crawler.engine._slot.scheduler.stop = True - - def parse(self, response): - pass - - actual_urls = [] - - def track_url(request, spider): - actual_urls.append(request.url) - - settings = {"SCHEDULER": TestScheduler, "SEEDING_POLICY": "lazy"} - crawler = get_crawler(TestSpider, settings_dict=settings) - crawler.signals.connect(track_url, signals.request_reached_downloader) - await maybe_deferred_to_future(crawler.crawl()) - assert crawler.stats.get_value("finish_reason") == "finished" - expected_urls = ["data:,a", "data:,b", "data:,c"] - assert actual_urls == expected_urls, f"{actual_urls=} != {expected_urls=}" - - @deferred_f_from_coro_f - async def test_lazy_start_order(self): - """By default, start requests should be sent in the order in which they - are iterated.""" - - class TestSpider(Spider): - name = "test" - start_urls = ["data:,a", "data:,b", "data:,c"] - - def parse(self, response): - pass - - actual_urls = [] - - def track_url(request, spider): - actual_urls.append(request.url) - - settings = {"SEEDING_POLICY": "lazy"} - crawler = get_crawler(TestSpider, settings_dict=settings) - crawler.signals.connect(track_url, signals.request_reached_downloader) - await maybe_deferred_to_future(crawler.crawl()) - assert crawler.stats.get_value("finish_reason") == "finished" - expected_urls = ["data:,a", "data:,b", "data:,c"] - assert actual_urls == expected_urls, f"{actual_urls=} != {expected_urls=}" - - @deferred_f_from_coro_f - async def test_front_load(self): - class TestScheduler(BaseScheduler): - def __init__(self, *args, **kwargs): - self.requests = defaultdict(deque) - - def enqueue_request(self, request: Request) -> bool: - self.requests[request.priority].append(request) - return True - - def has_pending_requests(self) -> bool: - return bool(self.requests) - - def next_request(self) -> Request | None: - if not self.requests: - return None - priority = max(self.requests) - request = self.requests[priority].popleft() - if not self.requests[priority]: - del self.requests[priority] - return request - - class TestSpider(Spider): - name = "test" - - async def start(self): - yield Request("data:,b", priority=0) - yield Request("data:,a", priority=1) - - def parse(self, response): - pass - - actual_urls = [] - - def track_url(request, spider): - actual_urls.append(request.url) - - settings = {"SCHEDULER": TestScheduler, "SEEDING_POLICY": "front_load"} - crawler = get_crawler(TestSpider, settings_dict=settings) - crawler.signals.connect(track_url, signals.request_reached_downloader) - await maybe_deferred_to_future(crawler.crawl()) - assert crawler.stats.get_value("finish_reason") == "finished" - expected_urls = ["data:,a", "data:,b"] - assert actual_urls == expected_urls, f"{actual_urls=} != {expected_urls=}" - - @deferred_f_from_coro_f - async def test_override(self): - class TestScheduler(BaseScheduler): - def __init__(self, *args, **kwargs): - self.requests = defaultdict(deque) - - def enqueue_request(self, request: Request) -> bool: - self.requests[request.priority].append(request) - return True - - def has_pending_requests(self) -> bool: - return bool(self.requests) - - def next_request(self) -> Request | None: - if not self.requests: - return None - priority = max(self.requests) - request = self.requests[priority].popleft() - if not self.requests[priority]: - del self.requests[priority] - return request - - class TestSpider(Spider): - name = "test" - - async def start(self): - yield "front-load" # typo - yield SeedingPolicy.front_load - yield Request("data:,b", priority=1) - yield Request("data:,a", priority=2) - yield self.crawler.settings["SEEDING_POLICY"] - yield Request("data:,c", priority=3) - - def parse(self, response): - pass - - actual_items = [] - actual_urls = [] - - def track_item(item, response, spider): - actual_items.append(item) - - def track_url(request, spider): - actual_urls.append(request.url) - - settings = {"SCHEDULER": TestScheduler, "SEEDING_POLICY": "lazy"} - crawler = get_crawler(TestSpider, settings_dict=settings) - crawler.signals.connect(track_item, signals.item_scraped) - crawler.signals.connect(track_url, signals.request_reached_downloader) - with LogCapture(level=ERROR) as log: - await maybe_deferred_to_future(crawler.crawl()) - assert len(log.records) == 1 - assert "must be valid seeding policies" in str(log.records[0]) - assert crawler.stats.get_value("finish_reason") == "finished" - assert not actual_items, ( - f"{actual_items=} should be empty, policies are not items" - ) - expected_urls = ["data:,a", "data:,b", "data:,c"] - assert actual_urls == expected_urls, f"{actual_urls=} != {expected_urls=}" - - -class MockServerTestCase(TestCase): - # See the comment on the matching line above. - timeout = ExecutionEngine._SLOT_HEARTBEAT_INTERVAL - # If requests are too fast, test_idle will fail because the outcome will - # match that of the lazy seeding policy. - delay = 0.2 - - @classmethod - def setUpClass(cls): - cls.mockserver = MockServer() - cls.mockserver.__enter__() - - @classmethod - def tearDownClass(cls): - cls.mockserver.__exit__(None, None, None) - - @deferred_f_from_coro_f - async def test_idle(self): - def _url(id): - return self.mockserver.url(f"/delay?n={self.delay}&{id}") - - class TestScheduler(BaseScheduler): - def __init__(self, *args, **kwargs): - self.requests = deque((Request(_url("a")),)) - - def enqueue_request(self, request: Request) -> bool: - self.requests.append(request) - return True - - def has_pending_requests(self) -> bool: - return bool(self.requests) - - def next_request(self) -> Request | None: - try: - return self.requests.popleft() - except IndexError: - return None - - class TestSpider(Spider): - name = "test" - start_urls = [_url("b"), _url("d")] - queue = deque((Request(_url("c")),)) - - def parse(self, response): - try: - request = self.queue.popleft() - except IndexError: - pass - else: - yield request - - actual_urls = [] - - def track_url(request, spider): - actual_urls.append(request.url) - - settings = {"SCHEDULER": TestScheduler, "SEEDING_POLICY": "idle"} - crawler = get_crawler(TestSpider, settings_dict=settings) - crawler.signals.connect(track_url, signals.request_reached_downloader) - await maybe_deferred_to_future(crawler.crawl()) - assert crawler.stats.get_value("finish_reason") == "finished" - expected_urls = [_url(letter) for letter in "abcd"] - assert actual_urls == expected_urls, f"{actual_urls=} != {expected_urls=}" diff --git a/tests/test_scheduler.py b/tests/test_scheduler.py index 1d6992a32..f90293dd3 100644 --- a/tests/test_scheduler.py +++ b/tests/test_scheduler.py @@ -3,6 +3,7 @@ from __future__ import annotations import shutil import tempfile from abc import ABC, abstractmethod +from collections import deque from typing import Any, NamedTuple import pytest @@ -10,7 +11,7 @@ from twisted.internet import defer from twisted.trial.unittest import TestCase from scrapy.core.downloader import Downloader -from scrapy.core.scheduler import Scheduler +from scrapy.core.scheduler import BaseScheduler, Scheduler from scrapy.crawler import Crawler from scrapy.http import Request from scrapy.spiders import Spider @@ -20,6 +21,38 @@ from scrapy.utils.test import get_crawler from tests.mockserver import MockServer +class MemoryScheduler(BaseScheduler): + paused = False + + def __init__(self, *args, **kwargs): + super().__init__(*args, **kwargs) + self.queue = deque( + Request(value) if isinstance(value, str) else value + for value in getattr(self, "queue", []) + ) + + def enqueue_request(self, request: Request) -> bool: + self.queue.append(request) + return True + + def has_pending_requests(self) -> bool: + return self.paused or bool(self.queue) + + def next_request(self) -> Request | None: + if self.paused: + return None + try: + return self.queue.pop() + except IndexError: + return None + + def pause(self) -> None: + self.paused = True + + def unpause(self) -> None: + self.paused = False + + class MockEngine(NamedTuple): downloader: MockDownloader