diff --git a/scrapy/core/downloader/__init__.py b/scrapy/core/downloader/__init__.py index 7ec86a08b..ca5e19983 100644 --- a/scrapy/core/downloader/__init__.py +++ b/scrapy/core/downloader/__init__.py @@ -118,7 +118,6 @@ class Downloader: ) self._slot_gc_loop: AsyncioLoopingCall | LoopingCall | None = None self._accepting_requests: bool = True - self._fast_stopping: bool = False self._download_tasks: dict[Request, Deferred[None]] = {} self.per_slot_settings: dict[str, dict[str, Any]] = self.settings.getdict( "DOWNLOAD_SLOTS" @@ -278,7 +277,6 @@ class Downloader: async def stop_async(self) -> int: self._accepting_requests = False - self._fast_stopping = True dropped_count = 0 diff --git a/tests/test_core_downloader.py b/tests/test_core_downloader.py index aba1541e0..5a3792e6f 100644 --- a/tests/test_core_downloader.py +++ b/tests/test_core_downloader.py @@ -2,11 +2,12 @@ from __future__ import annotations import warnings from typing import TYPE_CHECKING, cast +from unittest.mock import patch import OpenSSL.SSL import pytest from pytest_twisted import async_yield_fixture -from twisted.internet.defer import Deferred +from twisted.internet.defer import CancelledError, Deferred from twisted.web import server, static from twisted.web.client import Agent, BrowserLikePolicyForHTTPS, readBody from twisted.web.client import Response as TxResponse @@ -244,3 +245,81 @@ async def test_stop_async_rejects_new_requests() -> None: match="not accepting new requests", ): await downloader._enqueue_request(Request("https://example.com")) + + +@coroutine_test +async def test_wait_for_download_errbacks_queue_deferred_on_error() -> None: + crawler = get_crawler(DefaultSpider) + downloader = Downloader(crawler) + slot = Slot(concurrency=1, delay=0, randomize_delay=False) + + queue_dfd: Deferred = Deferred() + failures: list[Failure] = [] + queue_dfd.addErrback(failures.append) + + with patch.object(downloader, "_download", side_effect=RuntimeError("boom")): + await downloader._wait_for_download( + slot, + Request("https://example.com"), + queue_dfd, + ) + + assert len(failures) == 1 + assert failures[0].check(RuntimeError) + + +@coroutine_test +async def test_wait_for_download_keeps_called_queue_deferred_on_error() -> None: + crawler = get_crawler(DefaultSpider) + downloader = Downloader(crawler) + slot = Slot(concurrency=1, delay=0, randomize_delay=False) + + queue_dfd: Deferred = Deferred() + queue_dfd.callback(None) + + with patch.object(downloader, "_download", side_effect=RuntimeError("boom")): + await downloader._wait_for_download( + slot, + Request("https://example.com"), + queue_dfd, + ) + + assert queue_dfd.called + + +@coroutine_test +async def test_stop_async_skips_called_queued_deferred() -> None: + crawler = get_crawler(DefaultSpider) + crawler.spider = crawler._create_spider() + downloader = Downloader(crawler) + slot = Slot(concurrency=1, delay=0, randomize_delay=False) + downloader.slots["example.com"] = slot + + queue_dfd: Deferred = Deferred() + queue_dfd.callback(None) + slot.queue.append((Request("https://example.com"), queue_dfd)) + + dropped = await downloader.stop_async() + assert dropped == 1 + + +@coroutine_test +async def test_stop_async_cancels_pending_download_tasks() -> None: + crawler = get_crawler(DefaultSpider) + downloader = Downloader(crawler) + + done_dfd: Deferred = Deferred() + done_dfd.callback(None) + + pending_dfd: Deferred = Deferred() + failures: list[Failure] = [] + pending_dfd.addErrback(failures.append) + + downloader._download_tasks[Request("https://done.example")] = done_dfd + downloader._download_tasks[Request("https://pending.example")] = pending_dfd + + dropped = await downloader.stop_async() + + assert dropped == 1 + assert len(failures) == 1 + assert failures[0].check(CancelledError) diff --git a/tests/test_crawler.py b/tests/test_crawler.py index 75a02aadc..9cc9a9d5d 100644 --- a/tests/test_crawler.py +++ b/tests/test_crawler.py @@ -851,3 +851,34 @@ async def test_crawler_force_stop_uses_force_callback() -> None: crawler._set_force_stop_callback(force_stop_callback) await crawler.stop_async(mode="force") assert called + + +@coroutine_test +async def test_crawler_stop_async_without_engine_is_noop() -> None: + crawler = get_crawler(DefaultSpider) + crawler.crawling = True + + await crawler.stop_async(mode="graceful") + + assert crawler.crawling is False + + +@coroutine_test +async def test_crawler_stop_async_ignores_engine_not_running_runtime_error() -> None: + crawler = get_crawler(DefaultSpider) + crawler.crawling = True + + class DummyEngine: + running = True + called = False + + async def stop_async(self, *, mode: str = "graceful") -> None: + self.called = True + raise RuntimeError("Engine not running") + + dummy_engine = DummyEngine() + crawler.engine = dummy_engine # type: ignore[assignment] + + await crawler.stop_async(mode="graceful") + + assert dummy_engine.called is True diff --git a/tests/test_engine.py b/tests/test_engine.py index 5900695af..4c408d13c 100644 --- a/tests/test_engine.py +++ b/tests/test_engine.py @@ -8,7 +8,7 @@ from collections import defaultdict from dataclasses import dataclass from logging import DEBUG from typing import TYPE_CHECKING, Any, cast -from unittest.mock import Mock, call, patch +from unittest.mock import AsyncMock, Mock, call, patch from urllib.parse import urlparse import attr @@ -17,11 +17,12 @@ from itemadapter import ItemAdapter from pydispatch import dispatcher from testfixtures import LogCapture from twisted.internet import defer +from twisted.python.failure import Failure 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.exceptions import CloseSpider, DownloadCancelledError, IgnoreRequest from scrapy.http import Headers, Request, Response from scrapy.item import Field, Item from scrapy.linkextractors import LinkExtractor @@ -39,8 +40,6 @@ 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 @@ -451,6 +450,57 @@ class TestEngine(TestEngineBase): yield deferred_from_coro(e.start_async()) yield deferred_from_coro(e.stop_async()) + @coroutine_test + async def test_stop_async_force_mode_not_supported(self) -> None: + engine = ExecutionEngine(get_crawler(DefaultSpider), lambda _: None) + + with pytest.raises(ValueError, match="force stop mode is not supported"): + await engine.stop_async(mode="force") + + @coroutine_test + async def test_stop_async_not_running_raises(self) -> None: + engine = ExecutionEngine(get_crawler(DefaultSpider), lambda _: None) + + with pytest.raises(RuntimeError, match="Engine not running"): + await engine.stop_async() + + @coroutine_test + async def test_stop_async_reentrant_fast_waits_for_closewait(self) -> None: + engine = ExecutionEngine(get_crawler(DefaultSpider), lambda _: None) + engine.spider = Mock() + engine._stopping = True + engine._closewait = defer.Deferred() + + with patch.object( + engine, "close_spider_async", new_callable=AsyncMock + ) as close: + task = asyncio.create_task(engine.stop_async(mode="fast")) + await asyncio.sleep(0) + close.assert_called_once_with(reason="shutdown", mode="fast") + assert not task.done() + + assert engine._closewait + engine._closewait.callback(None) + await task + + @coroutine_test + async def test_handle_downloader_output_ignores_fast_cancelled_failures( + self, + ) -> None: + engine = ExecutionEngine(get_crawler(DefaultSpider), lambda _: None) + engine.spider = Mock() + engine._stop_mode = "fast" + + enqueue_scrape = Mock() + engine.scraper.enqueue_scrape = enqueue_scrape # type: ignore[method-assign] + + result = Failure(DownloadCancelledError("dropped during fast stop")) + await maybe_deferred_to_future( + engine._handle_downloader_output(result, Request("https://example.com")) + ) + + enqueue_scrape.assert_not_called() + @pytest.mark.only_asyncio @coroutine_test async def test_start_already_running_exception_asyncio(self): @@ -790,3 +840,23 @@ class TestEngineCloseSpider: assert calls == 1 assert crawler.stats assert crawler.stats.get_value("downloader/request_dropped_count") == 3 + + @coroutine_test + async def test_fast_stop_downloader_is_idempotent(self, crawler: Crawler) -> None: + engine = ExecutionEngine(crawler, lambda _: None) + crawler.engine = engine + await engine.open_spider_async() + + calls = 0 + + async def fast_stop_downloader() -> int: + nonlocal calls + calls += 1 + return 1 + + with patch.object(engine.downloader, "stop_async", fast_stop_downloader): + await engine._fast_stop_downloader() + await engine._fast_stop_downloader() + + assert calls == 1 + await engine.close_spider_async()