mirror of https://github.com/scrapy/scrapy.git
Cover the new error handling in news.rst and cover it with a test
This commit is contained in:
parent
0a137e40b7
commit
734f4e9a1a
|
|
@ -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
|
||||
<topics-jobs>`.
|
||||
|
||||
The logged errors have also been improved.
|
||||
|
||||
(:issue:`5426`)
|
||||
|
||||
Bug fixes
|
||||
~~~~~~~~~
|
||||
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
|
||||
|
|
|
|||
|
|
@ -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 = []
|
||||
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
Loading…
Reference in New Issue