diff --git a/docs/news.rst b/docs/news.rst index 79e7e1365..34d0b0880 100644 --- a/docs/news.rst +++ b/docs/news.rst @@ -90,6 +90,18 @@ New features (:issue:`3463`, :issue:`4058`, :issue:`6148`, :issue:`6715`, :issue:`6728`) +- Unlike its precedesors, if :meth:`~scrapy.Spider.yield_seeds` or + :meth:`~scrapy.spidermiddlewares.SpiderMiddleware.process_seeds` are not + asynchronous generators, the crawl starts nonetheless, without seeds. + + This aligns with the behavior of spider callbacks, and allows the crawl to + process requests from the scheduler, e.g. when :ref:`resuming a paused job + `. + + The logged errors have also been improved. + + (:issue:`5426`) + Bug fixes ~~~~~~~~~ diff --git a/tests/test_engine.py b/tests/test_engine.py index 223e89326..155060bf0 100644 --- a/tests/test_engine.py +++ b/tests/test_engine.py @@ -30,7 +30,6 @@ from twisted.web import server, static, util from scrapy import signals from scrapy.core.engine import ExecutionEngine, _Slot -from scrapy.core.scheduler import BaseScheduler from scrapy.exceptions import CloseSpider, IgnoreRequest from scrapy.http import Request from scrapy.item import Field, Item @@ -40,6 +39,7 @@ from scrapy.spiders import Spider from scrapy.utils.signal import disconnect_all from scrapy.utils.test import get_crawler from tests import get_testdata, tests_datadir +from tests.test_scheduler import MemoryScheduler class MyItem(Item): @@ -476,14 +476,6 @@ class TestEngine(TestEngineBase): def test_request_scheduled_signal(caplog): - class TestScheduler(BaseScheduler): - def __init__(self): - self.enqueued = [] - - def enqueue_request(self, request: Request) -> bool: - self.enqueued.append(request) - return True - def signal_handler(request: Request, spider: Spider) -> None: if "drop" in request.url: raise IgnoreRequest @@ -492,7 +484,7 @@ def test_request_scheduled_signal(caplog): crawler = get_crawler(spider.__class__) engine = ExecutionEngine(crawler, lambda _: None) engine.downloader._slot_gc_loop.stop() - scheduler = TestScheduler() + scheduler = MemoryScheduler() async def seeds(): return @@ -506,8 +498,8 @@ def test_request_scheduled_signal(caplog): drop_request = Request("https://drop.example") caplog.set_level(DEBUG) engine._schedule_request(drop_request, spider) - assert scheduler.enqueued == [keep_request], ( - f"{scheduler.enqueued!r} != [{keep_request!r}]" + assert list(scheduler.queue) == [keep_request], ( + f"{list(scheduler.queue)!r} != [{keep_request!r}]" ) crawler.signals.disconnect(signal_handler, request_scheduled) diff --git a/tests/test_engine_seeding.py b/tests/test_engine_seeding.py index c538fae44..aee7fad33 100644 --- a/tests/test_engine_seeding.py +++ b/tests/test_engine_seeding.py @@ -1,6 +1,6 @@ from __future__ import annotations -from collections import defaultdict, deque +from collections import deque from logging import ERROR from testfixtures import LogCapture @@ -8,12 +8,11 @@ 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 tests.test_scheduler import MemoryScheduler, PriorityScheduler from .mockserver import MockServer -from .test_spider_yield_seeds import twisted_sleep class MainTestCase(TestCase): @@ -33,22 +32,8 @@ class MainTestCase(TestCase): @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 TestScheduler(MemoryScheduler): + queue = ["data:,b"] class TestSpider(Spider): name = "test" @@ -72,22 +57,8 @@ class MainTestCase(TestCase): @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 TestScheduler(MemoryScheduler): + queue = ["data:,a"] class TestSpider(Spider): name = "test" @@ -114,35 +85,20 @@ class MainTestCase(TestCase): """If the scheduler reports having requests but yields none, the lazy policy schedules requests from seeds.""" - class TestScheduler(BaseScheduler): + class TestScheduler(MemoryScheduler): def __init__(self, *args, **kwargs): - self.requests = deque() + super().__init__(*args, **kwargs) 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 yield_seeds(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:,a") yield Request("data:,c") self.crawler.engine._slot.scheduler.stop = True @@ -157,10 +113,13 @@ class MainTestCase(TestCase): 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()) + with LogCapture() as log: + 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=}" + assert actual_urls == expected_urls, ( + f"{actual_urls=} != {expected_urls=}\n{log}" + ) @deferred_f_from_coro_f async def test_lazy_seed_order(self): @@ -189,26 +148,6 @@ class MainTestCase(TestCase): @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" @@ -224,7 +163,7 @@ class MainTestCase(TestCase): def track_url(request, spider): actual_urls.append(request.url) - settings = {"SCHEDULER": TestScheduler, "SEEDING_POLICY": "front_load"} + settings = {"SCHEDULER": PriorityScheduler, "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()) @@ -234,26 +173,6 @@ class MainTestCase(TestCase): @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" @@ -277,7 +196,7 @@ class MainTestCase(TestCase): def track_url(request, spider): actual_urls.append(request.url) - settings = {"SCHEDULER": TestScheduler, "SEEDING_POLICY": "lazy"} + settings = {"SCHEDULER": PriorityScheduler, "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) @@ -314,35 +233,21 @@ class MockServerTestCase(TestCase): 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 TestScheduler(MemoryScheduler): + queue = [_url("a")] class TestSpider(Spider): name = "test" start_urls = [_url("b"), _url("d")] - queue = deque((Request(_url("c")),)) + queue = deque([_url("c")]) def parse(self, response): try: - request = self.queue.popleft() + url = self.queue.popleft() except IndexError: pass else: - yield request + yield Request(url) actual_urls = [] diff --git a/tests/test_scheduler.py b/tests/test_scheduler.py index 1d6992a32..0c61c6e7d 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 defaultdict, 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,50 @@ from scrapy.utils.test import get_crawler from tests.mockserver import MockServer +class MemoryScheduler(BaseScheduler): + def __init__(self, *args, **kwargs): + super().__init__(*args, **kwargs) + self.queue = deque( + value if isinstance(value, Request) else Request(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 bool(self.queue) + + def next_request(self) -> Request | None: + try: + return self.queue.pop() + except IndexError: + return None + + +class PriorityScheduler(BaseScheduler): + def __init__(self, *args, **kwargs): + super().__init__(*args, **kwargs) + self.queue = defaultdict(deque) + + def enqueue_request(self, request: Request) -> bool: + self.queue[request.priority].append(request) + return True + + def has_pending_requests(self) -> bool: + return bool(self.queue) + + def next_request(self) -> Request | None: + if not self.queue: + return None + priority = max(self.queue) + request = self.queue[priority].popleft() + if not self.queue[priority]: + del self.queue[priority] + return request + + class MockEngine(NamedTuple): downloader: MockDownloader diff --git a/tests/test_spider_yield_seeds.py b/tests/test_spider_yield_seeds.py index d214ad45f..44030abd0 100644 --- a/tests/test_spider_yield_seeds.py +++ b/tests/test_spider_yield_seeds.py @@ -11,6 +11,8 @@ from scrapy.exceptions import ScrapyDeprecationWarning from scrapy.utils.defer import deferred_f_from_coro_f, maybe_deferred_to_future from scrapy.utils.test import get_crawler +from .test_scheduler import MemoryScheduler + # These are the minimum seconds necessary to wait to reproduce the issue that # has been solved by catching the RuntimeError exception in the # ExecutionEngine._next_request() method. A lower value makes these tests pass @@ -33,18 +35,19 @@ def twisted_sleep(seconds): class MainTestCase(TestCase): # Utility methods - async def _test_spider(self, spider, expected_items=None): + async def _test_spider(self, spider, expected_items=None, settings=None): actual_items = [] expected_items = [] if expected_items is None else expected_items + settings = settings or {} def track_item(item, response, spider): actual_items.append(item) - crawler = get_crawler(spider) + crawler = get_crawler(spider, settings_dict=settings) crawler.signals.connect(track_item, signals.item_scraped) await maybe_deferred_to_future(crawler.crawl()) assert crawler.stats.get_value("finish_reason") == "finished" - assert actual_items == expected_items + assert actual_items == expected_items, f"{actual_items=} != {expected_items=}" async def _test_yield_seeds(self, yield_seeds_, expected_items=None): class TestSpider(Spider): @@ -169,6 +172,31 @@ class MainTestCase(TestCase): assert ".yield_seeds must be an async generator function" in str(log) + @deferred_f_from_coro_f + async def test_bad_definition_continuance(self): + """Even if yield_seeds (or process_seeds) are not correctly defined, + blocking the iteration of seeds, requests from the scheduler are still + consumed.""" + + class TestScheduler(MemoryScheduler): + queue = ["data:,"] + + class TestSpider(Spider): + name = "test" + + async def yield_seeds(self): + return + + async def parse(self, response): + yield ITEM_A + + settings = {"SCHEDULER": TestScheduler} + + with LogCapture() as log: + await self._test_spider(TestSpider, [ITEM_A], settings=settings) + + assert ".yield_seeds must be an async generator function" in str(log), log + # Exceptions during iteration. @deferred_f_from_coro_f