From 003466496b1f7aecee0e97d56066beb31875039e Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Adri=C3=A1n=20Chaves?= Date: Mon, 10 Mar 2025 09:16:37 +0100 Subject: [PATCH] Fix a >5s asyncio.sleep() raising RuntimeError --- conftest.py | 15 ++++------ scrapy/core/engine.py | 16 ++++++---- tests/test_spider_yield_seeds.py | 51 ++++++++++++++++++++++++++++++++ tox.ini | 1 + 4 files changed, 68 insertions(+), 15 deletions(-) create mode 100644 tests/test_spider_yield_seeds.py diff --git a/conftest.py b/conftest.py index f33ffb1a4..106f4c9e0 100644 --- a/conftest.py +++ b/conftest.py @@ -48,14 +48,6 @@ def chdir(tmpdir): tmpdir.chdir() -def pytest_addoption(parser): - parser.addoption( - "--reactor", - default="default", - choices=["default", "asyncio"], - ) - - @pytest.fixture(scope="class") def reactor_pytest(request): if not request.cls: @@ -66,8 +58,11 @@ def reactor_pytest(request): @pytest.fixture(autouse=True) -def only_asyncio(request, reactor_pytest): - if request.node.get_closest_marker("only_asyncio") and reactor_pytest != "asyncio": +def only_asyncio(request): + if ( + request.node.get_closest_marker("only_asyncio") + and request.config.getoption("--reactor") != "asyncio" + ): pytest.skip("This test is only run with --reactor=asyncio") diff --git a/scrapy/core/engine.py b/scrapy/core/engine.py index 42fbb7219..14c7fb7b3 100644 --- a/scrapy/core/engine.py +++ b/scrapy/core/engine.py @@ -88,6 +88,8 @@ class _SeedingPolicy(Enum): class ExecutionEngine: + _SLOT_HEARTBEAT_INTERVAL: float = 5.0 + def __init__( self, crawler: Crawler, @@ -116,7 +118,7 @@ class ExecutionEngine: def _load_seeding_policy(self) -> None: try: - policy = _SeedingPolicy(self.settings["SEEDING_POLICY"]) + self._seeding_policy = _SeedingPolicy(self.settings["SEEDING_POLICY"]) except ValueError: supported_values = ", ".join(policy.value for policy in _SeedingPolicy) raise ValueError( @@ -124,7 +126,6 @@ class ExecutionEngine: f"({self.settings['SEEDING_POLICY']!r}) is not supported. " f"Supported values: {supported_values}." ) - self._feed = getattr(self, f"_{policy.name}_feed") def _get_scheduler_class(self, settings: BaseSettings) -> type[BaseScheduler]: from scrapy.core.scheduler import BaseScheduler @@ -187,7 +188,7 @@ class ExecutionEngine: self.paused = False @inlineCallbacks - def _lazy_feed(self) -> Generator[Deferred[Any], Any, None]: + def _next_request(self) -> Generator[Deferred[Any], Any, None]: if self.slot is None: return @@ -207,6 +208,11 @@ class ExecutionEngine: request_or_item = yield deferred_from_coro(self.slot.seeds.__anext__()) except StopAsyncIteration: self.slot.seeds = None + except RuntimeError: + # “RuntimeError: anext(): asynchronous generator is already + # running” happens if yield_seeds is taking long to yield the + # next seed. + pass except Exception: self.slot.seeds = None logger.error( @@ -395,7 +401,7 @@ class ExecutionEngine: 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}) - nextcall = CallLaterOnce(self._feed) + nextcall = CallLaterOnce(self._next_request) scheduler = build_from_crawler(self.scheduler_cls, self.crawler) seeds = yield self.scraper.spidermw.process_seeds(spider) self.slot = Slot(close_if_idle, nextcall, scheduler, seeds=seeds) @@ -407,7 +413,7 @@ class ExecutionEngine: self.crawler.stats.open_spider(spider) yield self.signals.send_catch_log_deferred(signals.spider_opened, spider=spider) self.slot.nextcall.schedule() - self.slot.heartbeat.start(5) + self.slot.heartbeat.start(self._SLOT_HEARTBEAT_INTERVAL) def _spider_idle(self) -> None: """ diff --git a/tests/test_spider_yield_seeds.py b/tests/test_spider_yield_seeds.py new file mode 100644 index 000000000..dcda67f29 --- /dev/null +++ b/tests/test_spider_yield_seeds.py @@ -0,0 +1,51 @@ +from asyncio import sleep + +import pytest +from pytest_twisted import ensureDeferred + +from scrapy import Spider, signals +from scrapy.core.engine import ExecutionEngine +from scrapy.utils.test import get_crawler + + +class Scenario: + pass + + +class AsyncioScenario(Scenario): + expected_items = [{"a": "b"}] + only_asyncio = True + + async def yield_seeds(self): + await sleep(ExecutionEngine._SLOT_HEARTBEAT_INTERVAL + 0.01) + yield {"a": "b"} + + +@pytest.mark.parametrize( + "scenario", + [ + pytest.param( + scenario, + marks=pytest.mark.only_asyncio + if getattr(scenario, "only_asyncio", False) + else [], + ) + for scenario in Scenario.__subclasses__() + ], +) +@ensureDeferred +async def test_main(scenario): + class TestSpider(Spider): + name = "test" + yield_seeds = scenario.yield_seeds + + actual_items = [] + + def track_item(item, response, spider): + actual_items.append(item) + + crawler = get_crawler(TestSpider) + crawler.signals.connect(track_item, signals.item_scraped) + await crawler.crawl() + assert crawler.stats.get_value("finish_reason") == "finished" + assert actual_items == scenario.expected_items diff --git a/tox.ini b/tox.ini index 041fcffca..71a58d52f 100644 --- a/tox.ini +++ b/tox.ini @@ -16,6 +16,7 @@ deps = pygments pytest != 8.2.* # https://github.com/pytest-dev/pytest/issues/12275 pytest-cov >= 4.0.0 + pytest-twisted pytest-xdist sybil >= 1.3.0 # https://github.com/cjw296/sybil/issues/20#issuecomment-605433422 testfixtures