from __future__ import annotations import asyncio import re import subprocess import sys from collections import defaultdict from dataclasses import dataclass from typing import TYPE_CHECKING, Any, cast from unittest.mock import Mock, call from urllib.parse import urlparse import attr import pytest from itemadapter import ItemAdapter from pydispatch import dispatcher from testfixtures import LogCapture from twisted.internet import defer 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 Headers, Request, Response from scrapy.item import Field, Item from scrapy.linkextractors import LinkExtractor from scrapy.spiders import Spider from scrapy.statscollectors import MemoryStatsCollector from scrapy.utils.defer import ( _schedule_coro, deferred_from_coro, maybe_deferred_to_future, ) from scrapy.utils.signal import disconnect_all from scrapy.utils.spider import DefaultSpider from scrapy.utils.test import get_crawler from tests import get_testdata from tests.utils.decorators import coroutine_test, inline_callbacks_test if TYPE_CHECKING: from twisted.python.failure import Failure from scrapy.core.scheduler import Scheduler from scrapy.crawler import Crawler from tests.mockserver.http import MockServer class MyItem(Item): name = Field() url = Field() price = Field() @attr.s class AttrsItem: name = attr.ib(default="") url = attr.ib(default="") price = attr.ib(default=0) @dataclass class DataClassItem: name: str = "" url: str = "" price: int = 0 class MySpider(Spider): name = "scrapytest.org" itemurl_re = re.compile(r"item\d+.html") name_re = re.compile(r"
File not found.
\n" b" \n" b"\n" ) elif run.getpath(request.url) == "/numbers": # signal was fired multiple times assert len(data) > 1 # bytes were received in order numbers = [str(x).encode("utf8") for x in range(2**18)] assert joined_data == b"".join(numbers) @staticmethod def _assert_signals_caught(run: CrawlerRun) -> None: assert signals.engine_started in run.signals_caught assert signals.engine_stopped in run.signals_caught assert signals.spider_opened in run.signals_caught assert signals.spider_idle in run.signals_caught assert signals.spider_closed in run.signals_caught assert signals.headers_received in run.signals_caught 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): @coroutine_test async def test_crawler(self, mockserver: MockServer) -> None: for spider in ( MySpider, DictItemsSpider, AttrsItemsSpider, DataClassItemsSpider, ): run = CrawlerRun(spider) await run.run(mockserver) self._assert_visited_urls(run) self._assert_scheduled_requests(run, count=9) self._assert_downloaded_responses(run, count=9) self._assert_scraped_items(run) self._assert_signals_caught(run) self._assert_headers_received(run) self._assert_bytes_received(run) @coroutine_test async def test_crawler_dupefilter(self, mockserver: MockServer) -> None: run = CrawlerRun(DupeFilterSpider) await run.run(mockserver) self._assert_scheduled_requests(run, count=8) self._assert_dropped_requests(run) @coroutine_test async def test_crawler_itemerror(self, mockserver: MockServer) -> None: run = CrawlerRun(ItemZeroDivisionErrorSpider) await run.run(mockserver) self._assert_items_error(run) @coroutine_test async def test_crawler_change_close_reason_on_idle( self, mockserver: MockServer ) -> None: run = CrawlerRun(ChangeCloseReasonSpider) await run.run(mockserver) assert { "spider": run.crawler.spider, "reason": "custom_reason", } == run.signals_caught[signals.spider_closed] @coroutine_test async def test_close_downloader(self): e = ExecutionEngine(get_crawler(MySpider), lambda _: None) await e.close_async() def test_close_without_downloader(self): class CustomException(Exception): pass class BadDownloader: def __init__(self, crawler): raise CustomException with pytest.raises(CustomException): ExecutionEngine( get_crawler(MySpider, {"DOWNLOADER": BadDownloader}), lambda _: None ) @inline_callbacks_test def test_start_already_running_exception(self): crawler = get_crawler(DefaultSpider) crawler.spider = crawler._create_spider() e = ExecutionEngine(crawler, lambda _: None) crawler.engine = e yield deferred_from_coro(e.open_spider_async()) _schedule_coro(e.start_async()) with pytest.raises(RuntimeError, match="Engine already running"): yield deferred_from_coro(e.start_async()) yield deferred_from_coro(e.stop_async()) @pytest.mark.only_asyncio @coroutine_test async def test_start_already_running_exception_asyncio(self): crawler = get_crawler(DefaultSpider) crawler.spider = crawler._create_spider() e = ExecutionEngine(crawler, lambda _: None) crawler.engine = e await e.open_spider_async() with pytest.raises(RuntimeError, match="Engine already running"): await asyncio.gather(e.start_async(), e.start_async()) await e.stop_async() @inline_callbacks_test def test_start_request_processing_exception(self): class BadRequestFingerprinter: def fingerprint(self, request): raise ValueError # to make Scheduler.enqueue_request() fail class SimpleSpider(Spider): name = "simple" async def start(self): yield Request("data:,") crawler = get_crawler( SimpleSpider, {"REQUEST_FINGERPRINTER_CLASS": BadRequestFingerprinter} ) with LogCapture() as log: yield crawler.crawl() assert "Error while processing requests from start()" in str(log) assert "Spider closed (shutdown)" in str(log) def test_short_timeout(self): args = ( sys.executable, "-m", "scrapy.cmdline", "fetch", "-s", "CLOSESPIDER_TIMEOUT=0.001", "-s", "LOG_LEVEL=DEBUG", "http://toscrape.com", ) p = subprocess.Popen( args, stdout=subprocess.DEVNULL, stderr=subprocess.PIPE, ) try: _, stderr = p.communicate(timeout=15) except subprocess.TimeoutExpired: p.kill() p.communicate() pytest.fail("Command took too much time to complete") stderr_str = stderr.decode("utf-8") assert "AttributeError" not in stderr_str, stderr_str assert "AssertionError" not in stderr_str, stderr_str class TestEngineDownloadAsync: """Test cases for ExecutionEngine.download_async().""" @pytest.fixture def engine(self) -> ExecutionEngine: crawler = get_crawler(MySpider) engine = ExecutionEngine(crawler, lambda _: None) engine.downloader.close() engine.downloader = Mock() engine._slot = Mock() engine._slot.inprogress = set() return engine @staticmethod async def _download(engine: ExecutionEngine, request: Request) -> Response: return await engine.download_async(request) @coroutine_test async def test_download_async_success(self, engine): """Test basic successful async download of a request.""" request = Request("http://example.com") response = Response("http://example.com", body=b"test body") engine.spider = Mock() engine.downloader.fetch.return_value = defer.succeed(response) engine._slot.add_request = Mock() engine._slot.remove_request = Mock() result = await self._download(engine, request) assert result == response engine._slot.add_request.assert_called_once_with(request) engine._slot.remove_request.assert_called_once_with(request) engine.downloader.fetch.assert_called_once_with(request) @coroutine_test async def test_download_async_redirect(self, engine): """Test async download with a redirect request.""" original_request = Request("http://example.com") redirect_request = Request("http://example.com/redirect") final_response = Response("http://example.com/redirect", body=b"redirected") # First call returns redirect request, second call returns final response engine.downloader.fetch.side_effect = [ defer.succeed(redirect_request), defer.succeed(final_response), ] engine.spider = Mock() engine._slot.add_request = Mock() engine._slot.remove_request = Mock() result = await self._download(engine, original_request) assert result == final_response assert engine.downloader.fetch.call_count == 2 engine._slot.add_request.assert_has_calls( [call(original_request), call(redirect_request)] ) engine._slot.remove_request.assert_has_calls( [call(original_request), call(redirect_request)] ) @coroutine_test async def test_download_async_no_spider(self, engine): """Test async download attempt when no spider is available.""" request = Request("http://example.com") engine.spider = None with pytest.raises(RuntimeError, match="No open spider to crawl:"): await self._download(engine, request) @coroutine_test async def test_download_async_failure(self, engine): """Test async download when the downloader raises an exception.""" request = Request("http://example.com") error = RuntimeError("Download failed") engine.spider = Mock() engine.downloader.fetch.return_value = defer.fail(error) engine._slot.add_request = Mock() engine._slot.remove_request = Mock() with pytest.raises(RuntimeError, match="Download failed"): await self._download(engine, request) engine._slot.add_request.assert_called_once_with(request) engine._slot.remove_request.assert_called_once_with(request) @coroutine_test async def test_download_async_fetch_needs_spider(self, engine): """A downloader whose fetch() requires a spider gets it passed in.""" engine._downloader_fetch_needs_spider = True request = Request("http://example.com") response = Response("http://example.com", body=b"test body") engine.spider = Mock() engine.downloader.fetch.return_value = defer.succeed(response) engine._slot.add_request = Mock() engine._slot.remove_request = Mock() result = await self._download(engine, request) assert result == response engine.downloader.fetch.assert_called_once_with(request, engine.spider) @pytest.mark.filterwarnings("ignore::scrapy.exceptions.ScrapyDeprecationWarning") class TestEngineDownload(TestEngineDownloadAsync): """Test cases for ExecutionEngine.download().""" @staticmethod async def _download(engine: ExecutionEngine, request: Request) -> Response: return await maybe_deferred_to_future(engine.download(request)) @coroutine_test async def test_request_scheduled_signal(): 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 crawler = get_crawler(MySpider) engine = ExecutionEngine(crawler, lambda _: None) scheduler = TestScheduler() async def start(): return yield engine._start = start() engine._slot = _Slot(False, Mock(), scheduler) crawler.signals.connect(signal_handler, signals.request_scheduled) keep_request = Request("https://keep.example") engine._schedule_request(keep_request) drop_request = Request("https://drop.example") engine._schedule_request(drop_request) assert scheduler.enqueued == [keep_request], ( f"{scheduler.enqueued!r} != [{keep_request!r}]" ) crawler.signals.disconnect(signal_handler, signals.request_scheduled) class TestEngineThrottling: @pytest.fixture def engine(self): crawler = get_crawler(MySpider) engine = ExecutionEngine(crawler, lambda _: None) yield engine engine.downloader.close() def test_pause_cancels_throttling_wakeup(self, engine): wakeup = Mock() engine._throttling_wakeup = wakeup engine.pause() assert engine.paused is True wakeup.cancel.assert_called_once_with() assert engine._throttling_wakeup is None engine.unpause() assert engine.paused is False @pytest.mark.requires_reactor # call_later() needs a reactor or asyncio loop def test_maybe_arm_throttling_wakeup_arms_timer(self, engine): scheduler = Mock() scheduler.has_pending_requests.return_value = True scheduler.get_next_request_delay.return_value = 5.0 engine._slot = Mock() engine._slot.scheduler = scheduler engine._maybe_arm_throttling_wakeup() assert engine._throttling_wakeup is not None # Cancel the scheduled reactor call so it does not leak into other tests. engine._cancel_throttling_wakeup() def test_maybe_arm_throttling_wakeup_no_delay(self, engine): scheduler = Mock() scheduler.has_pending_requests.return_value = True scheduler.get_next_request_delay.return_value = None engine._slot = Mock() engine._slot.scheduler = scheduler engine._maybe_arm_throttling_wakeup() assert engine._throttling_wakeup is None def test_maybe_arm_throttling_wakeup_zero_delay(self, engine): # A 0 delay means a request is ready but could not be sent (e.g. the # downloader is at capacity); arming a 0-second timer would busy-loop # the engine, so no timer must be armed. scheduler = Mock() scheduler.has_pending_requests.return_value = True scheduler.get_next_request_delay.return_value = 0.0 engine._slot = Mock() engine._slot.scheduler = scheduler engine._maybe_arm_throttling_wakeup() assert engine._throttling_wakeup is None def test_warn_delayed_requests(self, engine): engine._delayed_requests_warn_threshold = 1 engine._throttling_waiting = {Request("http://a.example")} engine._slot = Mock() # A scheduler without get_next_request_delay is not throttling-aware, so the # warning recommends switching to one. engine._slot.scheduler = Mock(spec=BaseScheduler) with LogCapture() as log: engine._maybe_warn_delayed_requests() # A second call is a no-op (the warning is emitted only once). engine._maybe_warn_delayed_requests() assert engine._delayed_requests_warned is True log_text = str(log) assert "requests held back by throttling" in log_text assert log_text.count("ThrottlingAwareScheduler") == 1 def test_warn_delayed_requests_throttling_aware(self, engine): engine._delayed_requests_warn_threshold = 1 engine._throttling_waiting = {Request("http://a.example")} engine._slot = Mock() # A throttling-aware scheduler (one with get_next_request_delay) holds # throttled requests itself, so no switch is recommended. engine._slot.scheduler = Mock() with LogCapture() as log: engine._maybe_warn_delayed_requests() assert "ThrottlingAwareScheduler" not in str(log) def test_spider_is_idle_false_while_scheduling(self, engine): engine._slot = Mock() engine.scraper.slot = Mock() engine.scraper.slot.is_idle.return_value = True engine.downloader = Mock() engine.downloader.active = [] engine._throttling_waiting = set() engine._start = None engine._scheduling = 1 # An in-flight async enqueue keeps the spider from being considered idle. assert engine.spider_is_idle() is False @coroutine_test async def test_enqueue_request_async_dropped(self, engine): scheduler = Mock() async def enqueue_request_async(request): return False scheduler.enqueue_request_async = enqueue_request_async engine._slot = Mock() engine._slot.scheduler = scheduler engine.spider = Mock() dropped = [] def on_dropped(request, spider): dropped.append(request) engine.signals.connect(on_dropped, signals.request_dropped, weak=False) engine._scheduling = 1 request = Request("http://a.example") await engine._enqueue_request_async(request) assert dropped == [request] assert engine._scheduling == 0 engine._slot.nextcall.schedule.assert_called_once_with() @coroutine_test async def test_enqueue_request_async_error(self, engine): scheduler = Mock() async def enqueue_request_async(request): raise RuntimeError("boom") scheduler.enqueue_request_async = enqueue_request_async engine._slot = Mock() engine._slot.scheduler = scheduler engine.spider = Mock() engine._scheduling = 1 with LogCapture() as log: await engine._enqueue_request_async(Request("http://a.example")) assert "Error while enqueuing request" in str(log) assert engine._scheduling == 0 engine._slot.nextcall.schedule.assert_called_once_with() @coroutine_test async def test_enqueue_request_async_slot_gone(self, engine): scheduler = Mock() async def enqueue_request_async(request): # The spider is closed while the enqueue is in flight. engine._slot = None return True scheduler.enqueue_request_async = enqueue_request_async slot = Mock() slot.scheduler = scheduler engine._slot = slot engine.spider = Mock() engine._scheduling = 1 await engine._enqueue_request_async(Request("http://a.example")) assert engine._scheduling == 0 # No reschedule is attempted once the slot is gone. slot.nextcall.schedule.assert_not_called() class TestEngineCloseSpider: """Tests for exception handling coverage during close_spider_async().""" @pytest.fixture def crawler(self) -> Crawler: crawler = get_crawler(DefaultSpider) crawler.spider = crawler._create_spider() return crawler @coroutine_test async def test_no_slot(self, crawler: Crawler) -> None: engine = ExecutionEngine(crawler, lambda _: None) crawler.engine = engine await engine.open_spider_async() slot = engine._slot engine._slot = None with pytest.raises(RuntimeError, match="Engine slot not assigned"): await engine.close_spider_async() # close it correctly engine._slot = slot await engine.close_spider_async() @coroutine_test async def test_no_spider(self, crawler: Crawler) -> None: engine = ExecutionEngine(crawler, lambda _: None) with pytest.raises(RuntimeError, match="Spider not opened"): await engine.close_spider_async() engine.downloader.close() # cleanup @coroutine_test async def test_exception_slot( self, crawler: Crawler, caplog: pytest.LogCaptureFixture ) -> None: engine = ExecutionEngine(crawler, lambda _: None) crawler.engine = engine await engine.open_spider_async() assert engine._slot del engine._slot.heartbeat await engine.close_spider_async() assert "Slot close failure" in caplog.text @coroutine_test async def test_exception_downloader( self, crawler: Crawler, caplog: pytest.LogCaptureFixture ) -> None: engine = ExecutionEngine(crawler, lambda _: None) crawler.engine = engine await engine.open_spider_async() engine.downloader.close = Mock(side_effect=Exception("close failed")) await engine.close_spider_async() assert "Downloader close failure" in caplog.text @coroutine_test async def test_exception_scraper( self, crawler: Crawler, caplog: pytest.LogCaptureFixture ) -> None: engine = ExecutionEngine(crawler, lambda _: None) crawler.engine = engine await engine.open_spider_async() engine.scraper.slot = None await engine.close_spider_async() assert "Scraper close failure" in caplog.text @coroutine_test async def test_exception_scheduler( self, crawler: Crawler, caplog: pytest.LogCaptureFixture ) -> None: engine = ExecutionEngine(crawler, lambda _: None) crawler.engine = engine await engine.open_spider_async() assert engine._slot del cast("Scheduler", engine._slot.scheduler).dqs await engine.close_spider_async() assert "Scheduler close failure" in caplog.text @coroutine_test async def test_exception_signal( self, crawler: Crawler, caplog: pytest.LogCaptureFixture ) -> None: engine = ExecutionEngine(crawler, lambda _: None) crawler.engine = engine await engine.open_spider_async() signal_manager = engine.signals del engine.signals await engine.close_spider_async() assert "Error while sending spider_close signal" in caplog.text # send the spider_closed signal to close various components await signal_manager.send_catch_log_async( signal=signals.spider_closed, spider=engine.spider, reason="cancelled", ) @coroutine_test async def test_exception_stats( self, crawler: Crawler, caplog: pytest.LogCaptureFixture ) -> None: engine = ExecutionEngine(crawler, lambda _: None) crawler.engine = engine await engine.open_spider_async() assert isinstance(crawler.stats, MemoryStatsCollector) del crawler.stats.spider_stats await engine.close_spider_async() assert "Stats close failure" in caplog.text @coroutine_test async def test_exception_callback( self, crawler: Crawler, caplog: pytest.LogCaptureFixture ) -> None: engine = ExecutionEngine(crawler, lambda _: defer.fail(ValueError())) crawler.engine = engine await engine.open_spider_async() await engine.close_spider_async() assert "Error running spider_closed_callback" in caplog.text @coroutine_test async def test_exception_async_callback( self, crawler: Crawler, caplog: pytest.LogCaptureFixture ) -> None: async def cb(_): raise ValueError engine = ExecutionEngine(crawler, cb) crawler.engine = engine await engine.open_spider_async() await engine.close_spider_async() assert "Error running spider_closed_callback" in caplog.text