From 313f9de28dc39acc5941d5a089930326c3981fdb Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Adri=C3=A1n=20Chaves?= Date: Wed, 26 Mar 2025 19:40:53 +0100 Subject: [PATCH] Implement lazy support --- docs/news.rst | 3 +- docs/topics/spiders.rst | 38 +- scrapy/commands/shell.py | 2 +- scrapy/core/engine.py | 77 ++-- scrapy/crawler.py | 5 +- scrapy/shell.py | 23 +- scrapy/signalmanager.py | 18 +- scrapy/utils/reactor.py | 19 +- .../__init__.py | 13 +- tests/test_downloadermiddleware.py | 44 ++- tests/test_engine.py | 24 +- tests/test_engine_loop.py | 350 +++--------------- 12 files changed, 236 insertions(+), 380 deletions(-) diff --git a/docs/news.rst b/docs/news.rst index eb7ecf8d5..27a06358f 100644 --- a/docs/news.rst +++ b/docs/news.rst @@ -22,7 +22,8 @@ Backward-incompatible changes As a result, the order in which start requests are sent may change. See :ref:`start-requests` for details and information on how to force start - request order. + request order or pause start request iteration while there are scheduled + requests. - In ``scrapy.core.engine.ExecutionEngine``: diff --git a/docs/topics/spiders.rst b/docs/topics/spiders.rst index 3d6216d60..c640609cc 100644 --- a/docs/topics/spiders.rst +++ b/docs/topics/spiders.rst @@ -373,16 +373,22 @@ Scrapy does not try to send :meth:`~scrapy.Spider.start` requests in order. Instead, it prioritizes reaching :setting:`CONCURRENT_REQUESTS` and :ref:`scheduling ` start requests. -To change that, override the :meth:`~scrapy.Spider.start` method to set -:attr:`Request.priority `. For example: +Forcing a start request order +----------------------------- + +To force a specific **request order**, override the +:meth:`~scrapy.Spider.start` method to set :attr:`Request.priority +`. For example: - To send start requests before other requests: .. code-block:: python async def start(self): - async for request in super().start(): - yield request.replace(priority=1) + async for item_or_request in super().start(): + if isinstance(item_or_request, Request): + item_or_request = item_or_request.replace(priority=1) + yield item_or_request - To send start requests in order: @@ -390,13 +396,33 @@ To change that, override the :meth:`~scrapy.Spider.start` method to set async def start(self): priority = len(self.start_urls) - async for request in super().start(): - yield request.replace(priority=priority) + async for item_or_request in super().start(): + if isinstance(item_or_request, Request): + item_or_request = item_or_request.replace(priority=priority) + yield item_or_request priority -= 1 You can also :ref:`customize the scheduler ` if you need more control over request prioritization. +Delaying start request iteration +-------------------------------- + +You can override the :meth:`~scrapy.Spider.start` method as follows to pause +its iteration whenever there are scheduled requests: + +.. code-block:: python + + async def start(self): + async for item_or_request in super().start(): + if self.crawler.engine.needs_backoff(): + await self.crawler.signals.wait_for(signals.scheduler_empty) + yield item_or_request + +This can help minimize the number of requests in the scheduler at any given +time, to minimize resource usage (memory or disk, depending on +:setting:`JOBDIR`). + .. _builtin-spiders: Generic Spiders diff --git a/scrapy/commands/shell.py b/scrapy/commands/shell.py index 4c3aca605..9dabfcd9c 100644 --- a/scrapy/commands/shell.py +++ b/scrapy/commands/shell.py @@ -85,7 +85,7 @@ class Command(ScrapyCommand): crawler._apply_settings() # The Shell class needs a persistent engine in the crawler crawler.engine = crawler._create_engine() - crawler.engine.start() + crawler.engine.start(_start_request_processing=False) self._start_crawler_thread() diff --git a/scrapy/core/engine.py b/scrapy/core/engine.py index e1f252ef2..1e8d01255 100644 --- a/scrapy/core/engine.py +++ b/scrapy/core/engine.py @@ -12,7 +12,7 @@ from time import time from traceback import format_exc from typing import TYPE_CHECKING, Any, TypeVar, cast -from twisted.internet.defer import Deferred, inlineCallbacks, succeed +from twisted.internet.defer import Deferred, succeed from twisted.internet.task import LoopingCall from twisted.python.failure import Failure @@ -20,13 +20,16 @@ from scrapy import signals from scrapy.core.scraper import Scraper, _HandleOutputDeferred from scrapy.exceptions import CloseSpider, DontCloseSpider, IgnoreRequest from scrapy.http import Request, Response -from scrapy.utils.defer import deferred_from_coro +from scrapy.utils.defer import ( + deferred_f_from_coro_f, + maybe_deferred_to_future, +) 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 if TYPE_CHECKING: - from collections.abc import AsyncIterable, Callable, Generator + from collections.abc import AsyncIterable, Callable from scrapy.core.downloader import Downloader from scrapy.core.scheduler import BaseScheduler @@ -99,6 +102,7 @@ class ExecutionEngine: ) self.start_time: float | None = None self._start: AsyncIterable[Any] | None = None + self._started_request_processing = False downloader_cls: type[Downloader] = load_object(self.settings["DOWNLOADER"]) try: self.scheduler_cls: type[BaseScheduler] = self._get_scheduler_class( @@ -121,22 +125,28 @@ class ExecutionEngine: ) return scheduler_cls - @inlineCallbacks - def start(self) -> Generator[Deferred[Any], Any, None]: + @deferred_f_from_coro_f + async def start(self, _start_request_processing=True) -> None: if self.running: raise RuntimeError("Engine already running") self.start_time = time() - yield self.signals.send_catch_log_deferred(signal=signals.engine_started) + await maybe_deferred_to_future( + self.signals.send_catch_log_deferred(signal=signals.engine_started) + ) self.running = True + if _start_request_processing: + self.start_request_processing() self._closewait: Deferred[None] = Deferred() - yield self._closewait + await maybe_deferred_to_future(self._closewait) def stop(self) -> Deferred[None]: """Gracefully stop the execution engine""" - @inlineCallbacks - def _finish_stopping_engine(_: Any) -> Generator[Deferred[Any], Any, None]: - yield self.signals.send_catch_log_deferred(signal=signals.engine_stopped) + @deferred_f_from_coro_f + async def _finish_stopping_engine(_: Any) -> None: + await maybe_deferred_to_future( + self.signals.send_catch_log_deferred(signal=signals.engine_stopped) + ) self._closewait.callback(None) if not self.running: @@ -170,10 +180,9 @@ class ExecutionEngine: def unpause(self) -> None: self.paused = False - @inlineCallbacks - def _process_next_spider_start_yield(self): + async def _process_next_spider_start_yield(self): try: - item_or_request = yield deferred_from_coro(self._start.__anext__()) + item_or_request = await self._start.__anext__() except StopAsyncIteration: self._start = None except Exception as exception: @@ -190,30 +199,37 @@ class ExecutionEngine: self.scraper.start_itemproc(item_or_request, response=None) self._slot.nextcall.schedule() - @inlineCallbacks - def _process_spider_start(self) -> Generator[Deferred[Any], Any, None]: + @deferred_f_from_coro_f + async def start_request_processing(self) -> None: """Process start items and requests in an asynchronous loop. Items are scraped. Requests are scheduled. """ + if self._started_request_processing: + raise RuntimeError("Request processing already started") + self._started_request_processing = True + assert self._slot is not None # typing + self._slot.nextcall.schedule() + self._slot.heartbeat.start(self._SLOT_HEARTBEAT_INTERVAL) while self._start is not None: - yield self._process_next_spider_start_yield() - if not self._needs_backout(): + await self._process_next_spider_start_yield() + if not self.needs_backout(): self._slot.nextcall.schedule() + await self._slot.nextcall.wait() def _start_next_requests(self) -> None: if self._slot is None or self._slot.closing is not None or self.paused: return - while not self._needs_backout(): + 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() - def _needs_backout(self) -> bool: + def needs_backout(self) -> bool: assert self._slot is not None # typing assert self.scraper.slot is not None # typing return ( @@ -377,29 +393,30 @@ class ExecutionEngine: dwld.addBoth(_on_complete) return dwld - @inlineCallbacks - def open_spider( + @deferred_f_from_coro_f + async def open_spider( self, spider: Spider, close_if_idle: bool = True, - ) -> Generator[Deferred[Any], Any, None]: + ) -> None: if self._slot is not None: raise RuntimeError(f"No free spider slot when opening {spider.name!r}") logger.info("Spider opened", extra={"spider": spider}) + self.spider = spider nextcall = CallLaterOnce(self._start_next_requests) scheduler = build_from_crawler(self.scheduler_cls, self.crawler) - self._start = yield self.scraper.spidermw.process_start(spider) self._slot = _Slot(close_if_idle, nextcall, scheduler) - self.spider = spider + self._start = await maybe_deferred_to_future( + self.scraper.spidermw.process_start(spider) + ) if hasattr(scheduler, "open") and (d := scheduler.open(spider)): - yield d - yield self.scraper.open_spider(spider) + await maybe_deferred_to_future(d) + await maybe_deferred_to_future(self.scraper.open_spider(spider)) 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) + await maybe_deferred_to_future( + self.signals.send_catch_log_deferred(signals.spider_opened, spider=spider) + ) def _spider_idle(self) -> None: """ diff --git a/scrapy/crawler.py b/scrapy/crawler.py index 88d618085..749096db5 100644 --- a/scrapy/crawler.py +++ b/scrapy/crawler.py @@ -136,6 +136,9 @@ class Crawler: "Overridden settings:\n%(settings)s", {"settings": pprint.pformat(d)} ) + # Cannot use @deferred_f_from_coro_f because that relies on the reactor + # being installed already, which is done within _apply_settings(), inside + # this method. @inlineCallbacks def crawl(self, *args: Any, **kwargs: Any) -> Generator[Deferred[Any], Any, None]: if self.crawling: @@ -152,7 +155,7 @@ class Crawler: self._update_root_log_handler() self.engine = self._create_engine() yield self.engine.open_spider(self.spider) - yield maybeDeferred(self.engine.start) + yield self.engine.start() except Exception: self.crawling = False if self.engine is not None: diff --git a/scrapy/shell.py b/scrapy/shell.py index 5e5e57a9a..14bc31137 100644 --- a/scrapy/shell.py +++ b/scrapy/shell.py @@ -24,6 +24,7 @@ from scrapy.spiders import Spider from scrapy.utils.conf import get_config from scrapy.utils.console import DEFAULT_PYTHON_SHELLS, start_python_console from scrapy.utils.datatypes import SequenceExclude +from scrapy.utils.defer import deferred_f_from_coro_f, maybe_deferred_to_future from scrapy.utils.misc import load_object from scrapy.utils.reactor import is_asyncio_reactor_installed, set_asyncio_event_loop from scrapy.utils.response import open_in_browser @@ -102,25 +103,33 @@ class Shell: # set the asyncio event loop for the current thread event_loop_path = self.crawler.settings["ASYNCIO_EVENT_LOOP"] set_asyncio_event_loop(event_loop_path) - spider = self._open_spider(request, spider) + + def crawl_request(_): + assert self.crawler.engine is not None + self.crawler.engine.crawl(request) + + d2 = self._open_spider(request, spider) + d2.addCallback(crawl_request) + d = _request_deferred(request) d.addCallback(lambda x: (x, spider)) - assert self.crawler.engine - self.crawler.engine.crawl(request) return d - def _open_spider(self, request: Request, spider: Spider | None) -> Spider: + @deferred_f_from_coro_f + async def _open_spider(self, request: Request, spider: Spider | None) -> None: if self.spider: - return self.spider + return if spider is None: spider = self.crawler.spider or self.crawler._create_spider() self.crawler.spider = spider assert self.crawler.engine - self.crawler.engine.open_spider(spider, close_if_idle=False) + await maybe_deferred_to_future( + self.crawler.engine.open_spider(spider, close_if_idle=False) + ) + self.crawler.engine.start_request_processing() self.spider = spider - return spider def fetch( self, diff --git a/scrapy/signalmanager.py b/scrapy/signalmanager.py index e106418d6..5bdc1b690 100644 --- a/scrapy/signalmanager.py +++ b/scrapy/signalmanager.py @@ -1,13 +1,12 @@ from __future__ import annotations -from typing import TYPE_CHECKING, Any +from typing import Any from pydispatch import dispatcher +from twisted.internet.defer import Deferred from scrapy.utils import signal as _signal - -if TYPE_CHECKING: - from twisted.internet.defer import Deferred +from scrapy.utils.defer import maybe_deferred_to_future class SignalManager: @@ -75,3 +74,14 @@ class SignalManager: """ kwargs.setdefault("sender", self.sender) _signal.disconnect_all(signal, **kwargs) + + async def wait_for(self, signal): + """Await the next *signal*.""" + d = Deferred() + + def handle(): + self.disconnect(handle, signal) + d.callback(None) + + self.connect(handle, signal) + await maybe_deferred_to_future(d) diff --git a/scrapy/utils/reactor.py b/scrapy/utils/reactor.py index 099c81f0e..184a95820 100644 --- a/scrapy/utils/reactor.py +++ b/scrapy/utils/reactor.py @@ -7,6 +7,7 @@ from typing import TYPE_CHECKING, Any, Generic, TypeVar from warnings import catch_warnings, filterwarnings from twisted.internet import asyncioreactor, error +from twisted.internet.defer import Deferred from scrapy.utils.misc import load_object @@ -54,6 +55,7 @@ class CallLaterOnce(Generic[_T]): self._a: tuple[Any, ...] = a self._kw: dict[str, Any] = kw self._call: DelayedCall | None = None + self._deferreds = [] def schedule(self, delay: float = 0) -> None: from twisted.internet import reactor @@ -66,8 +68,23 @@ class CallLaterOnce(Generic[_T]): self._call.cancel() def __call__(self) -> _T: + from twisted.internet import reactor + self._call = None - return self._func(*self._a, **self._kw) + result = self._func(*self._a, **self._kw) + + for d in self._deferreds: + reactor.callLater(0, d.callback, None) + self._deferreds = [] + + return result + + async def wait(self): + from scrapy.utils.defer import maybe_deferred_to_future + + d = Deferred() + self._deferreds.append(d) + await maybe_deferred_to_future(d) def set_asyncio_event_loop_policy() -> None: diff --git a/tests/test_cmdline_crawl_with_pipeline/__init__.py b/tests/test_cmdline_crawl_with_pipeline/__init__.py index 5228f6abd..954f2c924 100644 --- a/tests/test_cmdline_crawl_with_pipeline/__init__.py +++ b/tests/test_cmdline_crawl_with_pipeline/__init__.py @@ -8,11 +8,16 @@ class TestCmdlineCrawlPipeline: args = (sys.executable, "-m", "scrapy.cmdline", "crawl", spname) cwd = Path(__file__).resolve().parent proc = Popen(args, stdout=PIPE, stderr=PIPE, cwd=cwd) - proc.communicate() - return proc.returncode + _, stderr = proc.communicate() + return proc.returncode, stderr def test_open_spider_normally_in_pipeline(self): - assert self._execute("normal") == 0 + returncode, stderr = self._execute("normal") + assert returncode == 0 def test_exception_at_open_spider_in_pipeline(self): - assert self._execute("exception") == 1 + returncode, stderr = self._execute("exception") + assert ( + returncode == 0 + ) # An unhandled exception in a pipeline should not stop the crawl + assert b'RuntimeError("exception")' in stderr diff --git a/tests/test_downloadermiddleware.py b/tests/test_downloadermiddleware.py index 5d1161750..6c061b330 100644 --- a/tests/test_downloadermiddleware.py +++ b/tests/test_downloadermiddleware.py @@ -12,6 +12,11 @@ from scrapy.core.downloader.middleware import DownloaderMiddlewareManager from scrapy.exceptions import _InvalidOutput from scrapy.http import Request, Response from scrapy.spiders import Spider +from scrapy.utils.defer import ( + deferred_f_from_coro_f, + deferred_to_future, + maybe_deferred_to_future, +) from scrapy.utils.python import to_bytes from scrapy.utils.test import get_crawler, get_from_asyncio_queue @@ -29,7 +34,7 @@ class TestManagerBase(TestCase): def tearDown(self): return self.crawler.engine.close_spider(self.spider) - def _download(self, request, response=None): + async def _download(self, request, response=None): """Executes downloader mw manager's download method and returns the result (Request or Response) or raise exception in case of failure. @@ -44,7 +49,7 @@ class TestManagerBase(TestCase): # catch deferred result and return the value results = [] dfd.addBoth(results.append) - self._wait(dfd) + await maybe_deferred_to_future(dfd) ret = results[0] if isinstance(ret, Failure): ret.raiseException() @@ -54,13 +59,15 @@ class TestManagerBase(TestCase): class TestDefaults(TestManagerBase): """Tests default behavior with default settings""" - def test_request_response(self): + @deferred_f_from_coro_f + async def test_request_response(self): req = Request("http://example.com/index.html") resp = Response(req.url, status=200) - ret = self._download(req, resp) + ret = await self._download(req, resp) assert isinstance(ret, Response), "Non-response returned" - def test_3xx_and_invalid_gzipped_body_must_redirect(self): + @deferred_f_from_coro_f + async def test_3xx_and_invalid_gzipped_body_must_redirect(self): """Regression test for a failure when redirecting a compressed request. @@ -85,13 +92,14 @@ class TestDefaults(TestManagerBase): "Location": "http://example.com/login", }, ) - ret = self._download(request=req, response=resp) + ret = await self._download(request=req, response=resp) assert isinstance(ret, Request), f"Not redirected: {ret!r}" assert to_bytes(ret.url) == resp.headers["Location"], ( "Not redirected to location header" ) - def test_200_and_invalid_gzipped_body_must_fail(self): + @deferred_f_from_coro_f + async def test_200_and_invalid_gzipped_body_must_fail(self): req = Request("http://example.com") body = b"

You are being redirected

" resp = Response( @@ -106,13 +114,14 @@ class TestDefaults(TestManagerBase): }, ) with pytest.raises(BadGzipFile): - self._download(request=req, response=resp) + await self._download(request=req, response=resp) class TestResponseFromProcessRequest(TestManagerBase): """Tests middleware returning a response from process_request.""" - def test_download_func_not_called(self): + @deferred_f_from_coro_f + async def test_download_func_not_called(self): resp = Response("http://example.com/index.html") class ResponseMiddleware: @@ -126,7 +135,7 @@ class TestResponseFromProcessRequest(TestManagerBase): dfd = self.mwman.download(download_func, req, self.spider) results = [] dfd.addBoth(results.append) - self._wait(dfd) + await maybe_deferred_to_future(dfd) assert results[0] is resp assert not download_func.called @@ -195,7 +204,8 @@ class TestProcessExceptionInvalidOutput(TestManagerBase): class TestMiddlewareUsingDeferreds(TestManagerBase): """Middlewares using Deferreds should work""" - def test_deferred(self): + @deferred_f_from_coro_f + async def test_deferred(self): resp = Response("http://example.com/index.html") class DeferredMiddleware: @@ -214,7 +224,7 @@ class TestMiddlewareUsingDeferreds(TestManagerBase): dfd = self.mwman.download(download_func, req, self.spider) results = [] dfd.addBoth(results.append) - self._wait(dfd) + await maybe_deferred_to_future(dfd) assert results[0] is resp assert not download_func.called @@ -224,7 +234,8 @@ class TestMiddlewareUsingDeferreds(TestManagerBase): class TestMiddlewareUsingCoro(TestManagerBase): """Middlewares using asyncio coroutines should work""" - def test_asyncdef(self): + @deferred_f_from_coro_f + async def test_asyncdef(self): resp = Response("http://example.com/index.html") class CoroMiddleware: @@ -238,13 +249,14 @@ class TestMiddlewareUsingCoro(TestManagerBase): dfd = self.mwman.download(download_func, req, self.spider) results = [] dfd.addBoth(results.append) - self._wait(dfd) + await maybe_deferred_to_future(dfd) assert results[0] is resp assert not download_func.called @pytest.mark.only_asyncio - def test_asyncdef_asyncio(self): + @deferred_f_from_coro_f + async def test_asyncdef_asyncio(self): resp = Response("http://example.com/index.html") class CoroMiddleware: @@ -258,7 +270,7 @@ class TestMiddlewareUsingCoro(TestManagerBase): dfd = self.mwman.download(download_func, req, self.spider) results = [] dfd.addBoth(results.append) - self._wait(dfd) + await deferred_to_future(dfd) assert results[0] is resp assert not download_func.called diff --git a/tests/test_engine.py b/tests/test_engine.py index 4dc77e438..381cb5bcf 100644 --- a/tests/test_engine.py +++ b/tests/test_engine.py @@ -150,7 +150,6 @@ class CrawlerRun: """A class to run the crawler and keep track of events occurred""" def __init__(self, spider_class): - self.spider = None self.respplug = [] self.reqplug = [] self.reqdropped = [] @@ -191,7 +190,6 @@ class CrawlerRun: self.response_downloaded, signals.response_downloaded ) self.crawler.crawl(start_urls=start_urls) - self.spider = self.crawler.spider self.deferred = defer.Deferred() dispatcher.connect(self.stop, signals.engine_stopped) @@ -297,7 +295,7 @@ class TestEngineBase(unittest.TestCase): assert len(run.itemerror) == 2 for item, response, spider, failure in run.itemerror: assert failure.value.__class__ is ZeroDivisionError - assert spider == run.spider + assert spider == run.crawler.spider assert item["url"] == response.url if "item1.html" in item["url"]: @@ -378,11 +376,14 @@ class TestEngineBase(unittest.TestCase): assert signals.spider_closed in run.signals_caught assert signals.headers_received in run.signals_caught - assert {"spider": run.spider} == run.signals_caught[signals.spider_opened] - assert {"spider": run.spider} == run.signals_caught[signals.spider_idle] - assert {"spider": run.spider, "reason": "finished"} == run.signals_caught[ - signals.spider_closed + assert {"spider": run.crawler.spider} == run.signals_caught[ + signals.spider_opened ] + assert {"spider": run.crawler.spider} == run.signals_caught[signals.spider_idle] + assert { + "spider": run.crawler.spider, + "reason": "finished", + } == run.signals_caught[signals.spider_closed] class TestEngine(TestEngineBase): @@ -420,9 +421,10 @@ class TestEngine(TestEngineBase): def test_crawler_change_close_reason_on_idle(self): run = CrawlerRun(ChangeCloseReasonSpider) yield run.run() - assert {"spider": run.spider, "reason": "custom_reason"} == run.signals_caught[ - signals.spider_closed - ] + assert { + "spider": run.crawler.spider, + "reason": "custom_reason", + } == run.signals_caught[signals.spider_closed] @defer.inlineCallbacks def test_close_downloader(self): @@ -472,7 +474,7 @@ class TestEngine(TestEngineBase): finally: timer.cancel() - assert b"Traceback" not in stderr + assert b"Traceback" not in stderr, stderr def test_request_scheduled_signal(caplog): diff --git a/tests/test_engine_loop.py b/tests/test_engine_loop.py index 9b7a06305..b2dd81b7f 100644 --- a/tests/test_engine_loop.py +++ b/tests/test_engine_loop.py @@ -1,6 +1,5 @@ from collections import deque -import pytest from twisted.internet.defer import Deferred from twisted.trial.unittest import TestCase @@ -13,12 +12,12 @@ from .mockserver import MockServer from .test_scheduler import MemoryScheduler -def sleep(seconds: float = 0.001): +async def sleep(seconds: float = 0.001) -> None: from twisted.internet import reactor deferred: Deferred[None] = Deferred() reactor.callLater(seconds, deferred.callback, None) - return maybe_deferred_to_future(deferred) + await maybe_deferred_to_future(deferred) class MainTestCase(TestCase): @@ -87,9 +86,7 @@ class RequestSendOrderTestCase(TestCase): callback requests have the same priority. It is a very unintuitive behavior, documented as “undefined” so that we may - change it in the future without breaking the contract. - - For the asyncio reactor: + change it in the future without breaking the contract: 1. First, the first CONCURRENT_REQUESTS start requests are sent in order. @@ -105,16 +102,16 @@ class RequestSendOrderTestCase(TestCase): but only when there are not enough pending requests yielded from callbacks to reach the configured concurrency. - For the default Twisted reactor, step 1 sends the last CONCURRENT_REQUESTS - start requests in reverse order instead. - The reverse order is because the scheduler uses a LIFO queue by default (SCHEDULER_MEMORY_QUEUE, SCHEDULER_DISK_QUEUE). The order of the first few requests is unnaffected because they are sent as soon as they are - scheduled, and the last start requests sent before callback requests are - those that can be sent before the first callback requests are scheduled. + scheduled. The last start requests sent before callback requests are those + that can be sent before the first callback requests are scheduled. """ + # Error out if any tests relies on the heartbeat. + timeout = ExecutionEngine._SLOT_HEARTBEAT_INTERVAL + @classmethod def setUpClass(cls): cls.mockserver = MockServer() @@ -177,11 +174,8 @@ class RequestSendOrderTestCase(TestCase): expected_nums = sorted(start_nums + cb_nums) assert actual_nums == expected_nums, f"{actual_nums=} != {expected_nums=}" - # Asyncio reactor behavior - - @pytest.mark.only_asyncio @deferred_f_from_coro_f - async def test_ar_default(self): + async def test_default(self): await maybe_deferred_to_future( self._test_request_order( start_nums=[ @@ -215,9 +209,8 @@ class RequestSendOrderTestCase(TestCase): ) ) - @pytest.mark.only_asyncio @deferred_f_from_coro_f - async def test_ar_conc1(self): + async def test_conc1(self): await maybe_deferred_to_future( self._test_request_order( start_nums=[1, 4, 2], @@ -226,9 +219,8 @@ class RequestSendOrderTestCase(TestCase): ) ) - @pytest.mark.only_asyncio @deferred_f_from_coro_f - async def test_ar_conc2(self): + async def test_conc2(self): await maybe_deferred_to_future( self._test_request_order( start_nums=[1, 2, 6, 4, 3], @@ -237,9 +229,8 @@ class RequestSendOrderTestCase(TestCase): ) ) - @pytest.mark.only_asyncio @deferred_f_from_coro_f - async def test_ar_conc8(self): + async def test_conc8(self): await maybe_deferred_to_future( self._test_request_order( start_nums=[1, 2, 3, 4, 5, 6, 7, 8, 18, 16, 15, 14, 13, 12, 11, 10, 9], @@ -248,9 +239,8 @@ class RequestSendOrderTestCase(TestCase): ) ) - @pytest.mark.only_asyncio @deferred_f_from_coro_f - async def test_ar_conc16(self): + async def test_conc16(self): await maybe_deferred_to_future( self._test_request_order( start_nums=[ @@ -293,9 +283,8 @@ class RequestSendOrderTestCase(TestCase): ) ) - @pytest.mark.only_asyncio @deferred_f_from_coro_f - async def test_ar_conc3_ds2(self): + async def test_conc3_ds2(self): await maybe_deferred_to_future( self._test_request_order( start_nums=[1, 2, 3, 8, 6, 5, 4], @@ -307,9 +296,8 @@ class RequestSendOrderTestCase(TestCase): ) ) - @pytest.mark.only_asyncio @deferred_f_from_coro_f - async def test_ar_tconc3_dconc2(self): + async def test_tconc3_dconc2(self): await maybe_deferred_to_future( self._test_request_order( start_nums=[1, 2, 3, 7, 5, 4], @@ -321,9 +309,8 @@ class RequestSendOrderTestCase(TestCase): ) ) - @pytest.mark.only_asyncio @deferred_f_from_coro_f - async def test_ar_tconc5_dconc3(self): + async def test_tconc5_dconc3(self): await maybe_deferred_to_future( self._test_request_order( start_nums=[1, 2, 3, 4, 5, 10, 8, 7, 6], @@ -335,9 +322,8 @@ class RequestSendOrderTestCase(TestCase): ) ) - @pytest.mark.only_asyncio @deferred_f_from_coro_f - async def test_ar_tconc5_dconc2_ds3(self): + async def test_tconc5_dconc2_ds3(self): await maybe_deferred_to_future( self._test_request_order( start_nums=[1, 2, 3, 4, 5, 12, 10, 9, 8, 7, 6], @@ -350,9 +336,8 @@ class RequestSendOrderTestCase(TestCase): ) ) - @pytest.mark.only_asyncio @deferred_f_from_coro_f - async def test_ar_tconc5_dconc3_ds2(self): + async def test_tconc5_dconc3_ds2(self): await maybe_deferred_to_future( self._test_request_order( start_nums=[1, 2, 3, 4, 5, 12, 10, 9, 8, 7, 6], @@ -365,9 +350,8 @@ class RequestSendOrderTestCase(TestCase): ) ) - @pytest.mark.only_asyncio @deferred_f_from_coro_f - async def test_ar_tconc7_dconc2_ds3(self): + async def test_tconc7_dconc2_ds3(self): await maybe_deferred_to_future( self._test_request_order( start_nums=[1, 2, 3, 4, 5, 6, 7, 15, 13, 12, 11, 10, 9, 8], @@ -380,9 +364,8 @@ class RequestSendOrderTestCase(TestCase): ) ) - @pytest.mark.only_asyncio @deferred_f_from_coro_f - async def test_ar_tconc7_dconc3_ds2(self): + async def test_tconc7_dconc3_ds2(self): await maybe_deferred_to_future( self._test_request_order( start_nums=[1, 2, 3, 4, 5, 6, 7, 15, 13, 12, 11, 10, 9, 8], @@ -395,9 +378,8 @@ class RequestSendOrderTestCase(TestCase): ) ) - @pytest.mark.only_asyncio @deferred_f_from_coro_f - async def test_ar_fast(self): + async def test_fast(self): """Very fast responses may increase the number of start requests sent in reverse order before the first callback request.""" await maybe_deferred_to_future( @@ -409,240 +391,6 @@ class RequestSendOrderTestCase(TestCase): ) ) - # Default reactor behavior - - @pytest.mark.only_not_asyncio - @deferred_f_from_coro_f - async def test_dr_default(self): - await maybe_deferred_to_future( - self._test_request_order( - start_nums=[ - 26, - 24, - 23, - 22, - 21, - 20, - 19, - 18, - 17, - 16, - 15, - 14, - 13, - 12, - 11, - 10, - 9, - 8, - 7, - 6, - 5, - 4, - 3, - 2, - 1, - ], - cb_nums=[25], - ) - ) - - @pytest.mark.only_not_asyncio - @deferred_f_from_coro_f - async def test_dr_conc1(self): - await maybe_deferred_to_future( - self._test_request_order( - start_nums=[4, 2, 1], - cb_nums=[3], - settings={"CONCURRENT_REQUESTS": 1}, - ) - ) - - @pytest.mark.only_not_asyncio - @deferred_f_from_coro_f - async def test_dr_conc2(self): - await maybe_deferred_to_future( - self._test_request_order( - start_nums=[6, 4, 3, 2, 1], - cb_nums=[5], - settings={"CONCURRENT_REQUESTS": 2}, - ) - ) - - @pytest.mark.only_not_asyncio - @deferred_f_from_coro_f - async def test_dr_conc8(self): - await maybe_deferred_to_future( - self._test_request_order( - start_nums=[18, 16, 15, 14, 13, 12, 11, 10, 9, 8, 7, 6, 5, 4, 3, 2, 1], - cb_nums=[17], - settings={"CONCURRENT_REQUESTS": 8}, - ) - ) - - @pytest.mark.only_not_asyncio - @deferred_f_from_coro_f - async def test_dr_conc16(self): - await maybe_deferred_to_future( - self._test_request_order( - start_nums=[ - 34, - 32, - 31, - 30, - 29, - 28, - 27, - 26, - 25, - 24, - 23, - 22, - 21, - 20, - 19, - 18, - 17, - 16, - 15, - 14, - 13, - 12, - 11, - 10, - 9, - 8, - 7, - 6, - 5, - 4, - 3, - 2, - 1, - ], - cb_nums=[33], - settings={"CONCURRENT_REQUESTS_PER_DOMAIN": 16}, - ) - ) - - @pytest.mark.only_not_asyncio - @deferred_f_from_coro_f - async def test_dr_conc3_ds2(self): - await maybe_deferred_to_future( - self._test_request_order( - start_nums=[8, 6, 5, 4, 3, 2, 1], - cb_nums=[7], - settings={ - "CONCURRENT_REQUESTS": 3, - }, - download_slots=2, - ) - ) - - @pytest.mark.only_not_asyncio - @deferred_f_from_coro_f - async def test_dr_tconc3_dconc2(self): - await maybe_deferred_to_future( - self._test_request_order( - start_nums=[7, 5, 4, 3, 2, 1], - cb_nums=[6], - settings={ - "CONCURRENT_REQUESTS": 3, - "CONCURRENT_REQUESTS_PER_DOMAIN": 2, - }, - ) - ) - - @pytest.mark.only_not_asyncio - @deferred_f_from_coro_f - async def test_dr_tconc5_dconc3(self): - await maybe_deferred_to_future( - self._test_request_order( - start_nums=[10, 8, 7, 6, 5, 4, 3, 2, 1], - cb_nums=[9], - settings={ - "CONCURRENT_REQUESTS": 5, - "CONCURRENT_REQUESTS_PER_DOMAIN": 3, - }, - ) - ) - - @pytest.mark.only_not_asyncio - @deferred_f_from_coro_f - async def test_dr_tconc5_dconc2_ds3(self): - await maybe_deferred_to_future( - self._test_request_order( - start_nums=[12, 10, 9, 8, 7, 6, 5, 4, 3, 2, 1], - cb_nums=[11], - settings={ - "CONCURRENT_REQUESTS": 5, - "CONCURRENT_REQUESTS_PER_DOMAIN": 2, - }, - download_slots=3, - ) - ) - - @pytest.mark.only_not_asyncio - @deferred_f_from_coro_f - async def test_dr_tconc5_dconc3_ds2(self): - await maybe_deferred_to_future( - self._test_request_order( - start_nums=[12, 10, 9, 8, 7, 6, 5, 4, 3, 2, 1], - cb_nums=[11], - settings={ - "CONCURRENT_REQUESTS": 5, - "CONCURRENT_REQUESTS_PER_DOMAIN": 3, - }, - download_slots=2, - ) - ) - - @pytest.mark.only_not_asyncio - @deferred_f_from_coro_f - async def test_dr_tconc7_dconc2_ds3(self): - await maybe_deferred_to_future( - self._test_request_order( - start_nums=[15, 13, 12, 11, 10, 9, 8, 7, 6, 5, 4, 3, 2, 1], - cb_nums=[14], - settings={ - "CONCURRENT_REQUESTS": 7, - "CONCURRENT_REQUESTS_PER_DOMAIN": 2, - }, - download_slots=3, - ) - ) - - @pytest.mark.only_not_asyncio - @deferred_f_from_coro_f - async def test_dr_tconc7_dconc3_ds2(self): - await maybe_deferred_to_future( - self._test_request_order( - start_nums=[15, 13, 12, 11, 10, 9, 8, 7, 6, 5, 4, 3, 2, 1], - cb_nums=[14], - settings={ - "CONCURRENT_REQUESTS": 7, - "CONCURRENT_REQUESTS_PER_DOMAIN": 3, - }, - download_slots=2, - ) - ) - - @pytest.mark.only_not_asyncio - @deferred_f_from_coro_f - async def test_dr_fast(self): - """Very fast responses may increase the number of start requests sent - in reverse order before the first callback request.""" - await maybe_deferred_to_future( - self._test_request_order( - start_nums=[3, 2, 1], - cb_nums=[4], - settings={"CONCURRENT_REQUESTS": 1}, - response_seconds=self.fast_seconds, - ) - ) - - # Behavior shared by both reactors - @deferred_f_from_coro_f async def test_await(self): """Awaiting slow operations in Spider.start() may lower the number of @@ -672,9 +420,8 @@ class RequestSendOrderTestCase(TestCase): # Examples from the “Start requests” section of the documentation about # spiders. - @pytest.mark.only_asyncio @deferred_f_from_coro_f - async def test_ar_start_requests_first(self): + async def test_start_requests_first(self): start_nums = [1, 3, 2] cb_nums = [4] response_seconds = self.slow_seconds @@ -695,29 +442,6 @@ class RequestSendOrderTestCase(TestCase): ) ) - @pytest.mark.only_not_asyncio - @deferred_f_from_coro_f - async def test_dr_start_requests_first(self): - start_nums = [3, 2, 1] - cb_nums = [4] - response_seconds = self.slow_seconds - download_slots = 1 - - async def start(spider): - for num in start_nums: - request = self._request(num, response_seconds, download_slots) - yield request.replace(priority=1) - - await maybe_deferred_to_future( - self._test_request_order( - start_nums=start_nums, - cb_nums=cb_nums, - settings={"CONCURRENT_REQUESTS": 1}, - response_seconds=response_seconds, - start_fn=start, - ) - ) - @deferred_f_from_coro_f async def test_start_requests_first_sorted(self): start_nums = [1, 2, 3] @@ -741,3 +465,33 @@ class RequestSendOrderTestCase(TestCase): start_fn=start, ) ) + + @deferred_f_from_coro_f + async def test_lazy(self): + start_nums = [1, 2, 4] + cb_nums = [3] + response_seconds = self.slow_seconds + download_slots = 1 + + async def start(spider): + for num in start_nums: + if spider.crawler.engine.needs_backout(): + await spider.crawler.signals.wait_for(signals.scheduler_empty) + request = self._request(num, response_seconds, download_slots) + yield request + + await maybe_deferred_to_future( + self._test_request_order( + start_nums=start_nums, + cb_nums=cb_nums, + settings={ + "CONCURRENT_REQUESTS": 1, + # Without the lazy approach, using the FIFO queue would + # yield a different result, with start requests not being + # sorted. + "SCHEDULER_MEMORY_QUEUE": "scrapy.squeues.FifoMemoryQueue", + }, + response_seconds=response_seconds, + start_fn=start, + ) + )