mirror of https://github.com/scrapy/scrapy.git
Fix a >5s asyncio.sleep() raising RuntimeError
This commit is contained in:
parent
da5ebf0bd7
commit
003466496b
15
conftest.py
15
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")
|
||||
|
||||
|
||||
|
|
|
|||
|
|
@ -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:
|
||||
"""
|
||||
|
|
|
|||
|
|
@ -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
|
||||
Loading…
Reference in New Issue