diff --git a/.github/workflows/checks.yml b/.github/workflows/checks.yml index ed2388a59..1bd42832b 100644 --- a/.github/workflows/checks.yml +++ b/.github/workflows/checks.yml @@ -13,6 +13,7 @@ concurrency: jobs: checks: runs-on: ubuntu-latest + timeout-minutes: 30 strategy: fail-fast: false matrix: @@ -53,6 +54,7 @@ jobs: pre-commit: runs-on: ubuntu-latest + timeout-minutes: 15 steps: - uses: actions/checkout@v6 - uses: pre-commit/action@v3.0.1 diff --git a/.github/workflows/tests-macos.yml b/.github/workflows/tests-macos.yml index 0409b3ef2..70407b488 100644 --- a/.github/workflows/tests-macos.yml +++ b/.github/workflows/tests-macos.yml @@ -13,6 +13,7 @@ concurrency: jobs: tests: runs-on: macos-latest + timeout-minutes: 60 env: PYTEST_ADDOPTS: -n auto strategy: diff --git a/.github/workflows/tests-ubuntu.yml b/.github/workflows/tests-ubuntu.yml index f51dd9799..c4d63e22b 100644 --- a/.github/workflows/tests-ubuntu.yml +++ b/.github/workflows/tests-ubuntu.yml @@ -13,6 +13,7 @@ concurrency: jobs: tests: runs-on: ubuntu-latest + timeout-minutes: 60 env: PYTEST_ADDOPTS: -n auto strategy: diff --git a/.github/workflows/tests-windows.yml b/.github/workflows/tests-windows.yml index f413782bc..b210aa7ed 100644 --- a/.github/workflows/tests-windows.yml +++ b/.github/workflows/tests-windows.yml @@ -13,6 +13,7 @@ concurrency: jobs: tests: runs-on: windows-latest + timeout-minutes: 60 env: PYTEST_ADDOPTS: -n auto strategy: diff --git a/scrapy/core/downloader/__init__.py b/scrapy/core/downloader/__init__.py index 7c0ee0eec..1bbd4aefd 100644 --- a/scrapy/core/downloader/__init__.py +++ b/scrapy/core/downloader/__init__.py @@ -13,6 +13,7 @@ from twisted.python.failure import Failure from scrapy import Request, Spider, signals from scrapy.core.downloader.handlers import DownloadHandlers from scrapy.core.downloader.middleware import DownloaderMiddlewareManager +from scrapy.exceptions import DownloadCancelledError from scrapy.resolver import dnscache from scrapy.utils.asyncio import ( AsyncioLoopingCall, @@ -23,7 +24,6 @@ from scrapy.utils.asyncio import ( from scrapy.utils.decorators import _warn_spider_arg from scrapy.utils.defer import ( _defer_sleep_async, - _schedule_coro, deferred_from_coro, maybe_deferred_to_future, ) @@ -117,6 +117,8 @@ class Downloader: DownloaderMiddlewareManager.from_crawler(crawler) ) self._slot_gc_loop: AsyncioLoopingCall | LoopingCall | None = None + self._accepting_requests: bool = True + self._download_tasks: dict[Request, Deferred[None]] = {} self.per_slot_settings: dict[str, dict[str, Any]] = self.settings.getdict( "DOWNLOAD_SLOTS" ) @@ -176,6 +178,10 @@ class Downloader: # passed as download_func into self.middleware.download() in self.fetch() async def _enqueue_request(self, request: Request) -> Response: + if not self._accepting_requests: + raise DownloadCancelledError( + "The downloader is shutting down and not accepting new requests" + ) key, slot = self._get_slot(request) request.meta[self.DOWNLOAD_SLOT] = key slot.active.add(request) @@ -210,12 +216,21 @@ class Downloader: while slot.queue and slot.free_transfer_slots() > 0: slot.lastseen = now request, queue_dfd = slot.queue.popleft() - _schedule_coro(self._wait_for_download(slot, request, queue_dfd)) + download_dfd = deferred_from_coro( + self._wait_for_download(slot, request, queue_dfd) + ) + assert isinstance(download_dfd, Deferred) + self._download_tasks[request] = download_dfd + download_dfd.addBoth(self._download_task_done, request) # prevent burst if inter-request delays were configured if delay: self._process_queue(slot) break + def _download_task_done(self, result: Any, request: Request) -> Any: + self._download_tasks.pop(request, None) + return result + def _latercall(self, slot: Slot) -> None: slot.latercall = None self._process_queue(slot) @@ -257,11 +272,48 @@ class Downloader: try: response = await self._download(slot, request) except Exception: - queue_dfd.errback(Failure()) + if not queue_dfd.called: + queue_dfd.errback(Failure()) else: queue_dfd.callback(response) # awaited in _enqueue_request() + async def stop_async(self) -> int: + self._accepting_requests = False + + dropped_count = 0 + + for slot in self.slots.values(): + slot.close() + while slot.queue: + request, queue_dfd = slot.queue.popleft() + dropped_count += 1 + if not queue_dfd.called: + queue_dfd.errback( + Failure( + DownloadCancelledError( + "Request dropped due to fast downloader shutdown" + ) + ) + ) + self.signals.send_catch_log( + signal=signals.request_left_downloader, + request=request, + spider=self.crawler.spider, + ) + + for download_dfd in list(self._download_tasks.values()): + if download_dfd.called: + continue + dropped_count += 1 + download_dfd.cancel() + + if dropped_count: + await _defer_sleep_async() + + return dropped_count + def close(self) -> None: + self._accepting_requests = False self._stop_slot_gc() for slot in self.slots.values(): slot.close() diff --git a/scrapy/core/engine.py b/scrapy/core/engine.py index 1033e874f..47672299d 100644 --- a/scrapy/core/engine.py +++ b/scrapy/core/engine.py @@ -25,10 +25,12 @@ from scrapy.core.scraper import Scraper from scrapy.exceptions import ( CloseSpider, DontCloseSpider, + DownloadCancelledError, IgnoreRequest, ScrapyDeprecationWarning, ) from scrapy.http import Request, Response +from scrapy.utils._stopmode import _max_stop_mode, _normalize_stop_mode, _StopMode from scrapy.utils.asyncio import ( AsyncioLoopingCall, create_looping_call, @@ -119,6 +121,8 @@ class ExecutionEngine: self.running: bool = False self._starting: bool = False self._stopping: bool = False + self._stop_mode: _StopMode = "graceful" + self._downloader_fast_stopped: bool = False self.paused: bool = False self._spider_closed_callback: Callable[ [Spider], Coroutine[Any, Any, None] | Deferred[None] | None @@ -202,23 +206,36 @@ class ExecutionEngine: with contextlib.suppress(asyncio.exceptions.CancelledError): await maybe_deferred_to_future(self._closewait) - def stop(self) -> Deferred[None]: # pragma: no cover + def stop( + self, *, mode: _StopMode = "graceful" + ) -> Deferred[None]: # pragma: no cover warnings.warn( "ExecutionEngine.stop() is deprecated, use stop_async() instead", ScrapyDeprecationWarning, stacklevel=2, ) - return deferred_from_coro(self.stop_async()) + return deferred_from_coro(self.stop_async(mode=mode)) - async def stop_async(self) -> None: + async def stop_async(self, *, mode: _StopMode = "graceful") -> None: """Gracefully stop the execution engine. .. versionadded:: 2.14 """ - if not self._starting: + mode = _normalize_stop_mode(mode, allow_force=False) + + if not self._starting and not self._stopping: raise RuntimeError("Engine not running") + self._stop_mode = _max_stop_mode(self._stop_mode, mode) + + if self._stopping: + if self.spider is not None and self._stop_mode == "fast": + await self.close_spider_async(reason="shutdown", mode="fast") + if self._closewait: + await maybe_deferred_to_future(self._closewait) + return + self.running = self._starting = False self._stopping = True if self._start_request_processing_awaitable is not None: @@ -233,7 +250,7 @@ class ExecutionEngine: self._start_request_processing_awaitable.cancel() self._start_request_processing_awaitable = None if self.spider is not None: - await self.close_spider_async(reason="shutdown") + await self.close_spider_async(reason="shutdown", mode=self._stop_mode) await self.signals.send_catch_log_async(signal=signals.engine_stopped) if self._closewait: self._closewait.callback(None) @@ -404,6 +421,17 @@ class ExecutionEngine: f"Incorrect type: expected Request, Response or Failure, got {type(result)}: {result!r}" ) + if ( + isinstance(result, Failure) + and self._stop_mode == "fast" + and result.check( + DownloadCancelledError, + CancelledError, + asyncio.exceptions.CancelledError, + ) + ): + return + # downloader middleware can return requests (for example, redirects) if isinstance(result, Request): self.crawl(result) @@ -583,20 +611,51 @@ class ExecutionEngine: _schedule_coro(self.close_spider_async(reason=ex.reason)) def close_spider( - self, spider: Spider, reason: str = "cancelled" + self, + spider: Spider, + reason: str = "cancelled", + mode: _StopMode = "graceful", ) -> Deferred[None]: # pragma: no cover warnings.warn( "ExecutionEngine.close_spider() is deprecated, use close_spider_async() instead", ScrapyDeprecationWarning, stacklevel=2, ) - return deferred_from_coro(self.close_spider_async(reason=reason)) + return deferred_from_coro(self.close_spider_async(reason=reason, mode=mode)) - async def close_spider_async(self, *, reason: str = "cancelled") -> None: # noqa: PLR0912 + async def _fast_stop_downloader(self) -> None: + if self._downloader_fast_stopped: + return + + self._downloader_fast_stopped = True + dropped_count = await self.downloader.stop_async() + + assert self.crawler.stats + if dropped_count: + self.crawler.stats.inc_value( + "downloader/request_dropped_count", dropped_count + ) + + logger.info( + "Fast shutdown dropped %(count)d downloader requests", + {"count": dropped_count}, + extra={"spider": self.spider}, + ) + + # pylint: disable=too-many-statements + async def close_spider_async( # noqa: PLR0912, PLR0915 + self, + *, + reason: str = "cancelled", + mode: _StopMode = "graceful", + ) -> None: """Close (cancel) spider and clear all its outstanding requests. .. versionadded:: 2.14 """ + mode = _normalize_stop_mode(mode, allow_force=False) + self._stop_mode = _max_stop_mode(self._stop_mode, mode) + if self.spider is None: raise RuntimeError("Spider not opened") @@ -604,6 +663,8 @@ class ExecutionEngine: raise RuntimeError("Engine slot not assigned") if self._slot.closing is not None: + if self._stop_mode == "fast": + await self._fast_stop_downloader() await maybe_deferred_to_future(self._slot.closing) return @@ -613,6 +674,9 @@ class ExecutionEngine: "Closing spider (%(reason)s)", {"reason": reason}, extra={"spider": spider} ) + if self._stop_mode == "fast": + await self._fast_stop_downloader() + try: await self._slot.close() except Exception: diff --git a/scrapy/crawler.py b/scrapy/crawler.py index 9828fe6f5..dd6a2b2f1 100644 --- a/scrapy/crawler.py +++ b/scrapy/crawler.py @@ -20,7 +20,8 @@ from scrapy.extension import ExtensionManager from scrapy.settings import SETTINGS_PRIORITIES, Settings, overridden_settings from scrapy.signalmanager import SignalManager from scrapy.spiderloader import SpiderLoaderProtocol, get_spider_loader -from scrapy.utils.defer import deferred_from_coro +from scrapy.utils._stopmode import _normalize_stop_mode, _StopMode +from scrapy.utils.defer import deferred_from_coro, ensure_awaitable from scrapy.utils.log import ( configure_logging, get_scrapy_root_handler, @@ -42,7 +43,7 @@ from scrapy.utils.reactor import ( from scrapy.utils.reactorless import install_reactor_import_hook if TYPE_CHECKING: - from collections.abc import Awaitable, Generator, Iterable + from collections.abc import Awaitable, Callable, Generator, Iterable from scrapy.logformatter import LogFormatter from scrapy.statscollectors import StatsCollector @@ -85,6 +86,15 @@ class Crawler: self.request_fingerprinter: RequestFingerprinterProtocol | None = None self.spider: Spider | None = None self.engine: ExecutionEngine | None = None + self._force_stop_callback: ( + Callable[[], Awaitable[None] | Deferred[None] | None] | None + ) = None + + def _set_force_stop_callback( + self, + callback: Callable[[], Awaitable[None] | Deferred[None] | None] | None, + ) -> None: + self._force_stop_callback = callback def _update_root_log_handler(self) -> None: if get_scrapy_root_handler() is not None: @@ -233,7 +243,7 @@ class Crawler: def _create_engine(self) -> ExecutionEngine: return ExecutionEngine(self, lambda _: self.stop_async()) - def stop(self) -> Deferred[None]: + def stop(self, *, mode: _StopMode = "graceful") -> Deferred[None]: """Start a graceful stop of the crawler and return a deferred that is fired when the crawler is stopped.""" warnings.warn( @@ -241,18 +251,42 @@ class Crawler: ScrapyDeprecationWarning, stacklevel=2, ) - return deferred_from_coro(self.stop_async()) + return deferred_from_coro(self.stop_async(mode=mode)) - async def stop_async(self) -> None: + async def stop_async(self, *, mode: _StopMode = "graceful") -> None: """Start a graceful stop of the crawler and complete when the crawler is stopped. .. versionadded:: 2.14 """ - if self.crawling: - self.crawling = False - assert self.engine - if self.engine.running: - await self.engine.stop_async() + mode = _normalize_stop_mode(mode) + was_crawling = self.crawling + self.crawling = False + + if not was_crawling and mode == "graceful": + return + + if mode == "force": + if self._force_stop_callback is not None: + await ensure_awaitable(self._force_stop_callback()) + return + logger.warning( + "Force stop requested, but no process-level force stop callback is available. Falling back to fast stop." + ) + mode = "fast" + + if self.engine is None: + return + + # During shutdown callbacks, graceful stop may be re-entered after + # the engine has already switched to non-running state. + if mode == "graceful" and not self.engine.running: + return + + try: + await self.engine.stop_async(mode=mode) + except RuntimeError as exc: + if str(exc) != "Engine not running": + raise @staticmethod def _get_component( @@ -475,13 +509,16 @@ class CrawlerRunner(CrawlerRunnerBase): self._active.discard(d) self.bootstrap_failed |= not getattr(crawler, "spider", None) or failed - def stop(self) -> Deferred[Any]: + def stop(self, *, mode: _StopMode = "graceful") -> Deferred[Any]: """ Stops simultaneously all the crawling jobs taking place. Returns a deferred that is fired when they all have ended. """ - return DeferredList(deferred_from_coro(c.stop_async()) for c in self.crawlers) + mode = _normalize_stop_mode(mode) + return DeferredList( + deferred_from_coro(c.stop_async(mode=mode)) for c in self.crawlers + ) @inlineCallbacks def join(self) -> Generator[Deferred[Any], Any, None]: @@ -599,15 +636,16 @@ class AsyncCrawlerRunner(CrawlerRunnerBase): return task - async def stop(self) -> None: + async def stop(self, *, mode: _StopMode = "graceful") -> None: """ Stops simultaneously all the crawling jobs taking place. Completes when they all have ended. """ + mode = _normalize_stop_mode(mode) if self.crawlers: await asyncio.wait( - [asyncio.create_task(c.stop_async()) for c in self.crawlers] + [asyncio.create_task(c.stop_async(mode=mode)) for c in self.crawlers] ) async def join(self) -> None: @@ -629,6 +667,14 @@ class CrawlerProcessBase(CrawlerRunnerBase): configure_logging(self.settings, install_root_handler) log_scrapy_info(self.settings) + def _create_crawler(self, spidercls: str | type[Spider]) -> Crawler: + crawler = super()._create_crawler(spidercls) + crawler._set_force_stop_callback(self._force_stop) + return crawler + + def _force_stop(self) -> None: + self._stop_reactor() + @abstractmethod def start( self, stop_after_crawl: bool = True, install_signal_handlers: bool = True @@ -638,10 +684,17 @@ class CrawlerProcessBase(CrawlerRunnerBase): def _signal_shutdown(self, signum: int, _: Any) -> None: from twisted.internet import reactor - install_shutdown_handlers(self._signal_kill) + install_shutdown_handlers(self._signal_fast_shutdown) self._log_shutdown(signum) reactor.callFromThread(self._graceful_stop_reactor) + def _signal_fast_shutdown(self, signum: int, _: Any) -> None: + from twisted.internet import reactor + + install_shutdown_handlers(self._signal_kill) + self._log_fast_shutdown(signum) + reactor.callFromThread(self._fast_stop_reactor) + def _signal_kill(self, signum: int, _: Any) -> None: from twisted.internet import reactor @@ -653,7 +706,15 @@ class CrawlerProcessBase(CrawlerRunnerBase): def _log_shutdown(signum: int) -> None: signame = signal_names[signum] logger.info( - "Received %(signame)s, shutting down gracefully. Send again to force ", + "Received %(signame)s, shutting down gracefully. Send again to stop faster.", + {"signame": signame}, + ) + + @staticmethod + def _log_fast_shutdown(signum: int) -> None: + signame = signal_names[signum] + logger.info( + "Received %(signame)s twice, dropping downloader requests. Send again to force unclean shutdown", {"signame": signame}, ) @@ -661,7 +722,8 @@ class CrawlerProcessBase(CrawlerRunnerBase): def _log_kill(signum: int) -> None: signame = signal_names[signum] logger.info( - "Received %(signame)s twice, forcing unclean shutdown", {"signame": signame} + "Received %(signame)s three times, forcing unclean shutdown", + {"signame": signame}, ) def _setup_reactor(self, install_signal_handlers: bool) -> None: @@ -703,13 +765,20 @@ class CrawlerProcessBase(CrawlerRunnerBase): ) @abstractmethod - def _stop_dfd(self) -> Deferred[Any]: + def _stop_dfd(self, *, mode: _StopMode = "graceful") -> Deferred[Any]: raise NotImplementedError @inlineCallbacks def _graceful_stop_reactor(self) -> Generator[Deferred[Any], Any, None]: try: - yield self._stop_dfd() + yield self._stop_dfd(mode="graceful") + finally: + self._stop_reactor() + + @inlineCallbacks + def _fast_stop_reactor(self) -> Generator[Deferred[Any], Any, None]: + try: + yield self._stop_dfd(mode="fast") finally: self._stop_reactor() @@ -762,10 +831,12 @@ class CrawlerProcess(CrawlerProcessBase, CrawlerRunner): spidercls = self.spider_loader.load(spidercls) init_reactor = not self._initialized_reactor self._initialized_reactor = True - return Crawler(spidercls, self.settings, init_reactor=init_reactor) + crawler = Crawler(spidercls, self.settings, init_reactor=init_reactor) + crawler._set_force_stop_callback(self._force_stop) + return crawler - def _stop_dfd(self) -> Deferred[Any]: - return self.stop() + def _stop_dfd(self, *, mode: _StopMode = "graceful") -> Deferred[Any]: + return self.stop(mode=mode) def start( self, stop_after_crawl: bool = True, install_signal_handlers: bool = True @@ -862,8 +933,18 @@ class AsyncCrawlerProcess(CrawlerProcessBase, AsyncCrawlerRunner): self._initialized_reactor = True self._reactorless_main_task: asyncio.Future[None] | None = None - def _stop_dfd(self) -> Deferred[Any]: - return deferred_from_coro(self.stop()) + def _force_stop(self) -> None: + if self.settings.getbool("TWISTED_REACTOR_ENABLED"): + self._stop_reactor() + return + + if (loop := self._reactorless_loop) is None: + return + if (task := self._reactorless_main_task) is not None: + loop.call_soon_threadsafe(task.cancel) + + def _stop_dfd(self, *, mode: _StopMode = "graceful") -> Deferred[Any]: + return deferred_from_coro(self.stop(mode=mode)) def start( self, stop_after_crawl: bool = True, install_signal_handlers: bool = True @@ -1008,24 +1089,31 @@ class AsyncCrawlerProcess(CrawlerProcessBase, AsyncCrawlerRunner): } ) - def _signal_shutdown_reactorless(self, signum: int, _: Any) -> None: - install_shutdown_handlers(self._signal_kill_reactorless) - self._log_shutdown(signum) + def _schedule_reactorless_shutdown(self, *, mode: _StopMode) -> None: if (loop := self._reactorless_loop) is None: return - loop.call_soon_threadsafe(self._create_shutdown_task) + def _create_shutdown_task() -> None: + coro = self._shutdown_reactorless(mode=mode) + try: + loop.create_task(coro) + except RuntimeError: + coro.close() - def _create_shutdown_task(self) -> None: - assert self._reactorless_loop - coro = self._shutdown_graceful_reactorless() - try: - self._reactorless_loop.create_task(coro) - except RuntimeError: - coro.close() + loop.call_soon_threadsafe(_create_shutdown_task) - async def _shutdown_graceful_reactorless(self) -> None: - await self.stop() + def _signal_shutdown_reactorless(self, signum: int, _: Any) -> None: + install_shutdown_handlers(self._signal_fast_shutdown_reactorless) + self._log_shutdown(signum) + self._schedule_reactorless_shutdown(mode="graceful") + + def _signal_fast_shutdown_reactorless(self, signum: int, _: Any) -> None: + install_shutdown_handlers(self._signal_kill_reactorless) + self._log_fast_shutdown(signum) + self._schedule_reactorless_shutdown(mode="fast") + + async def _shutdown_reactorless(self, *, mode: _StopMode) -> None: + await self.stop(mode=mode) if not self._stop_after_crawl: # wait until crawl tasks finish and cancel the future await self.join() diff --git a/scrapy/utils/_stopmode.py b/scrapy/utils/_stopmode.py new file mode 100644 index 000000000..ec4d083a6 --- /dev/null +++ b/scrapy/utils/_stopmode.py @@ -0,0 +1,27 @@ +from __future__ import annotations + +from typing import Literal + +_StopMode = Literal["graceful", "fast", "force"] + +_STOP_MODE_PRIORITY: dict[_StopMode, int] = { + "graceful": 0, + "fast": 1, + "force": 2, +} + + +def _normalize_stop_mode(mode: _StopMode, *, allow_force: bool = True) -> _StopMode: + if mode not in _STOP_MODE_PRIORITY: + raise ValueError( + f"Unknown stop mode {mode!r}. Expected one of: graceful, fast, force" + ) + if mode == "force" and not allow_force: + raise ValueError("The force stop mode is not supported in this context") + return mode + + +def _max_stop_mode(mode1: _StopMode, mode2: _StopMode) -> _StopMode: + if _STOP_MODE_PRIORITY[mode1] >= _STOP_MODE_PRIORITY[mode2]: + return mode1 + return mode2 diff --git a/tests/AsyncCrawlerProcess/reactorless_sleeping_no_stop_after_crawl.py b/tests/AsyncCrawlerProcess/reactorless_sleeping_no_stop_after_crawl.py new file mode 100644 index 000000000..d3eb2bec1 --- /dev/null +++ b/tests/AsyncCrawlerProcess/reactorless_sleeping_no_stop_after_crawl.py @@ -0,0 +1,20 @@ +import asyncio +import sys + +import scrapy +from scrapy.crawler import AsyncCrawlerProcess + + +class SleepingSpider(scrapy.Spider): + name = "sleeping" + + start_urls = ["data:,;"] + + async def parse(self, response): + await asyncio.sleep(int(sys.argv[1])) + + +process = AsyncCrawlerProcess(settings={"TWISTED_REACTOR_ENABLED": False}) + +process.crawl(SleepingSpider) +process.start(stop_after_crawl=False) diff --git a/tests/test_core_downloader.py b/tests/test_core_downloader.py index fdd5edc27..37dfacfdc 100644 --- a/tests/test_core_downloader.py +++ b/tests/test_core_downloader.py @@ -1,11 +1,13 @@ from __future__ import annotations import warnings -from typing import TYPE_CHECKING, cast +from typing import TYPE_CHECKING, Any, cast +from unittest.mock import patch import OpenSSL.SSL import pytest from pytest_twisted import async_yield_fixture +from twisted.internet.defer import CancelledError, Deferred from twisted.internet.protocol import Factory from twisted.internet.protocol import Protocol as TxProtocol from twisted.internet.ssl import AcceptableCiphers, optionsForClientTLS @@ -14,13 +16,14 @@ from twisted.web import server, static from twisted.web.client import Agent, BrowserLikePolicyForHTTPS, readBody from twisted.web.client import Response as TxResponse +from scrapy import Request from scrapy.core.downloader import Downloader, Slot, tls from scrapy.core.downloader.contextfactory import ( _load_context_factory_from_settings, _ScrapyClientContextFactory, ) from scrapy.core.downloader.handlers.http11 import _RequestBodyProducer -from scrapy.exceptions import ScrapyDeprecationWarning +from scrapy.exceptions import DownloadCancelledError, ScrapyDeprecationWarning from scrapy.utils._deps_compat import ( PYOPENSSL_SET_CIPHER_LIST_TMP_CONN, TWISTED_TLS_NEW_IMPL, @@ -35,7 +38,7 @@ from tests.mockserver.utils import ssl_context_factory from tests.utils.decorators import coroutine_test if TYPE_CHECKING: - from twisted.internet.defer import Deferred + from twisted.python.failure import Failure from twisted.web.iweb import IBodyProducer @@ -310,6 +313,119 @@ async def test_fetch_deprecated_spider_arg(): await crawler.crawl_async() +@coroutine_test +async def test_stop_async_drops_queued_requests() -> 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 + + request = Request("https://example.com") + queue_dfd: Deferred[Any] = Deferred() + failures: list[Failure] = [] + queue_dfd.addErrback(failures.append) + slot.queue.append((request, queue_dfd)) + + dropped = await downloader.stop_async() + assert dropped == 1 + assert len(failures) == 1 + assert failures[0].check(DownloadCancelledError) + + +@coroutine_test +async def test_stop_async_rejects_new_requests() -> None: + crawler = get_crawler(DefaultSpider) + crawler.spider = crawler._create_spider() + downloader = Downloader(crawler) + + await downloader.stop_async() + + with pytest.raises( + DownloadCancelledError, + 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[Any] = 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[Any] = 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[Any] = 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[None] = Deferred() + done_dfd.callback(None) + + pending_dfd: Deferred[None] = 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) + + def test_deprecated_tls_module_names() -> None: with pytest.warns( ScrapyDeprecationWarning, diff --git a/tests/test_crawler.py b/tests/test_crawler.py index 0cddfd0ed..8a609e0fa 100644 --- a/tests/test_crawler.py +++ b/tests/test_crawler.py @@ -4,7 +4,8 @@ import asyncio import logging import re from pathlib import Path -from typing import Any, ClassVar +from typing import Any, ClassVar, cast +from unittest.mock import patch import pytest from zope.interface.exceptions import MultipleInvalid @@ -826,3 +827,282 @@ async def test_deprecated_crawler_stop() -> None: ScrapyDeprecationWarning, match=r"Crawler.stop\(\) is deprecated" ): await maybe_deferred_to_future(crawler.stop()) + + +@coroutine_test +async def test_crawler_stop_async_invalid_mode() -> None: + crawler = get_crawler(DefaultSpider) + with pytest.raises(ValueError, match=r"Unknown stop mode"): + await crawler.stop_async(mode="invalid") # type: ignore[arg-type] + + +@coroutine_test +async def test_crawler_graceful_stop_non_running_engine_is_noop() -> None: + crawler = get_crawler(DefaultSpider) + crawler.crawling = True + + class DummyEngine: + running = False + called = False + + async def stop_async(self, *, mode: str = "graceful") -> None: + self.called = True + + dummy_engine = DummyEngine() + crawler.engine = dummy_engine # type: ignore[assignment] + + await crawler.stop_async(mode="graceful") + + assert dummy_engine.called is False + + +@coroutine_test +async def test_crawler_force_stop_falls_back_to_fast( + caplog: pytest.LogCaptureFixture, +) -> None: + crawler = get_crawler(DefaultSpider) + + class DummyEngine: + called_mode: str | None = None + + async def stop_async(self, *, mode: str = "graceful") -> None: + self.called_mode = mode + + dummy_engine = DummyEngine() + crawler.engine = dummy_engine # type: ignore[assignment] + + with caplog.at_level(logging.WARNING): + await crawler.stop_async(mode="force") + + assert dummy_engine.called_mode == "fast" + assert "Falling back to fast stop" in caplog.text + + +@coroutine_test +async def test_crawler_force_stop_uses_force_callback() -> None: + crawler = get_crawler(DefaultSpider) + called = False + + def force_stop_callback() -> None: + nonlocal called + called = True + + 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 + + +@coroutine_test +async def test_crawler_stop_async_reraises_other_runtime_errors() -> None: + crawler = get_crawler(DefaultSpider) + crawler.crawling = True + + class DummyEngine: + running = True + + async def stop_async(self, *, mode: str = "graceful") -> None: + raise RuntimeError("different runtime error") + + crawler.engine = DummyEngine() # type: ignore[assignment] + + with pytest.raises(RuntimeError, match="different runtime error"): + await crawler.stop_async(mode="graceful") + + +@pytest.mark.requires_reactor +@coroutine_test +async def test_crawler_process_force_stop_via_public_crawler_api() -> None: + crawler_process = CrawlerProcess(install_root_handler=False) + called = False + + def stop_reactor(_: Any = None) -> None: + nonlocal called + called = True + + crawler_process._stop_reactor = stop_reactor # type: ignore[method-assign] + crawler = crawler_process.create_crawler(DefaultSpider) + + await crawler.stop_async(mode="force") + + assert called + + +@pytest.mark.only_asyncio +@pytest.mark.requires_reactor +@coroutine_test +async def test_async_crawler_process_force_stop_reactor_enabled_via_public_crawler_api() -> ( + None +): + crawler_process = AsyncCrawlerProcess( + {"TWISTED_REACTOR_ENABLED": True}, + install_root_handler=False, + ) + called = False + + def stop_reactor(_: Any = None) -> None: + nonlocal called + called = True + + crawler_process._stop_reactor = stop_reactor # type: ignore[method-assign] + crawler = crawler_process.create_crawler(DefaultSpider) + + await crawler.stop_async(mode="force") + + assert called + + +@pytest.mark.only_asyncio +@coroutine_test +async def test_async_crawler_process_force_stop_reactorless_without_main_task( + reactor_pytest: str, +) -> None: + if reactor_pytest != "none": + pytest.skip("This test is only for --reactor=none") + + crawler_process = AsyncCrawlerProcess( + {"TWISTED_REACTOR_ENABLED": False}, + install_root_handler=False, + ) + assert crawler_process._reactorless_loop is not None + assert crawler_process._reactorless_main_task is None + crawler = crawler_process.create_crawler(DefaultSpider) + + await crawler.stop_async(mode="force") + + +@pytest.mark.only_asyncio +@coroutine_test +async def test_async_crawler_process_force_stop_reactorless_without_loop( + reactor_pytest: str, +) -> None: + if reactor_pytest != "none": + pytest.skip("This test is only for --reactor=none") + + crawler_process = AsyncCrawlerProcess( + {"TWISTED_REACTOR_ENABLED": False}, + install_root_handler=False, + ) + crawler_process._reactorless_loop = None + crawler_process._reactorless_main_task = None + crawler = crawler_process.create_crawler(DefaultSpider) + + await crawler.stop_async(mode="force") + + +@pytest.mark.only_asyncio +@coroutine_test +async def test_async_crawler_process_force_stop_reactorless_with_task( + reactor_pytest: str, +) -> None: + if reactor_pytest != "none": + pytest.skip("This test is only for --reactor=none") + + crawler_process = AsyncCrawlerProcess( + {"TWISTED_REACTOR_ENABLED": False}, + install_root_handler=False, + ) + + class DummyLoop: + callback = None + + def call_soon_threadsafe(self, callback) -> None: + self.callback = callback + + class DummyTask: + called = False + + def cancel(self) -> None: + self.called = True + + loop = DummyLoop() + task = DummyTask() + crawler_process._reactorless_loop = cast("asyncio.AbstractEventLoop", loop) + crawler_process._reactorless_main_task = cast("asyncio.Future[None]", task) + crawler = crawler_process.create_crawler(DefaultSpider) + + await crawler.stop_async(mode="force") + + assert loop.callback is not None + loop.callback() + assert task.called is True + + +def test_async_crawler_process_schedule_reactorless_shutdown_without_loop() -> None: + crawler_process = object.__new__(AsyncCrawlerProcess) + crawler_process._reactorless_loop = None + + crawler_process._schedule_reactorless_shutdown(mode="graceful") + + +def test_async_crawler_process_schedule_reactorless_shutdown_runtime_error() -> None: + crawler_process = object.__new__(AsyncCrawlerProcess) + + class DummyLoop: + scheduled = False + create_task_called = False + + def call_soon_threadsafe(self, callback) -> None: + self.scheduled = True + callback() + + def create_task(self, coro) -> None: + self.create_task_called = True + raise RuntimeError("event loop is closing") + + class DummyCoro: + closed = False + + def close(self) -> None: + self.closed = True + + loop = DummyLoop() + coro = DummyCoro() + called_mode: str | None = None + + def shutdown_reactorless(*, mode: str) -> DummyCoro: + nonlocal called_mode + called_mode = mode + return coro + + crawler_process._reactorless_loop = cast("asyncio.AbstractEventLoop", loop) + with patch.object( + crawler_process, "_shutdown_reactorless", new=shutdown_reactorless + ): + crawler_process._schedule_reactorless_shutdown(mode="graceful") + + assert called_mode == "graceful" + + assert loop.scheduled + assert loop.create_task_called + assert coro.closed diff --git a/tests/test_crawler_subprocess.py b/tests/test_crawler_subprocess.py index e3f9ac161..7d8b44db0 100644 --- a/tests/test_crawler_subprocess.py +++ b/tests/test_crawler_subprocess.py @@ -231,10 +231,13 @@ class TestCrawlerProcessSubprocessBase(ScriptRunnerMixin): p.expect_exact("Crawled (200)") p.kill(sig) p.expect_exact("shutting down gracefully") - # sending the second signal too fast often causes problems + # sending a new signal too fast often causes problems await async_sleep(0.01) p.kill(sig) - p.expect_exact("forcing unclean shutdown") + p.expect_exact("dropping downloader requests") + await async_sleep(0.01) + p.kill(sig) + p.expect_exact("forcing unclean shutdown", timeout=20) p.wait() # type: ignore[no-untyped-call] if p.proc.stdin: p.proc.stdin.close() @@ -425,6 +428,9 @@ class TestAsyncCrawlerProcessSubprocess(TestCrawlerProcessSubprocessBase): def test_shutdown_graceful(self) -> None: self._test_shutdown_graceful("reactorless_sleeping.py") + def test_shutdown_graceful_stop_after_crawl_false(self) -> None: + self._test_shutdown_graceful("reactorless_sleeping_no_stop_after_crawl.py") + @coroutine_test async def test_shutdown_forced(self) -> None: await self._test_shutdown_forced("reactorless_sleeping.py") diff --git a/tests/test_engine.py b/tests/test_engine.py index e51eb4664..4dbe6a073 100644 --- a/tests/test_engine.py +++ b/tests/test_engine.py @@ -7,7 +7,7 @@ import sys from collections import defaultdict from dataclasses import dataclass from typing import TYPE_CHECKING, Any, cast -from unittest.mock import Mock, call +from unittest.mock import AsyncMock, Mock, call, patch from urllib.parse import urlparse import attr @@ -16,11 +16,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 @@ -35,11 +36,10 @@ 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 import async_sleep 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 +451,71 @@ 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: + stop_dfd = deferred_from_coro(engine.stop_async(mode="fast")) + await async_sleep(0) + close.assert_called_once_with(reason="shutdown", mode="fast") + assert not stop_dfd.called + + assert engine._closewait + engine._closewait.callback(None) + await maybe_deferred_to_future(stop_dfd) + + @coroutine_test + async def test_stop_async_reentrant_graceful_without_spider_or_closewait( + self, + ) -> None: + engine = ExecutionEngine(get_crawler(DefaultSpider), lambda _: None) + engine._stopping = True + + with patch.object( + engine, "close_spider_async", new_callable=AsyncMock + ) as close: + await engine.stop_async(mode="graceful") + + close.assert_not_called() + + @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): @@ -767,3 +832,45 @@ class TestEngineCloseSpider: await engine.open_spider_async() await engine.close_spider_async() assert "Error running spider_closed_callback" in caplog.text + + @coroutine_test + async def test_fast_close_stops_downloader_and_records_dropped_requests( + 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 3 + + with patch.object(engine.downloader, "stop_async", fast_stop_downloader): + await engine.close_spider_async(mode="fast") + + 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()