From 8613d390d7a2264fa8556d52e481607c0ac98242 Mon Sep 17 00:00:00 2001 From: Adrian Chaves Date: Thu, 23 Apr 2026 13:46:10 +0200 Subject: [PATCH 01/18] Support fast crawler stops --- scrapy/core/downloader/__init__.py | 60 +++++++++++- scrapy/core/engine.py | 79 ++++++++++++++-- scrapy/crawler.py | 146 +++++++++++++++++++++++------ scrapy/utils/_stopmode.py | 30 ++++++ tests/test_core_downloader.py | 40 +++++++- tests/test_crawler.py | 42 +++++++++ tests/test_crawler_subprocess.py | 5 +- tests/test_engine.py | 23 +++++ 8 files changed, 381 insertions(+), 44 deletions(-) create mode 100644 scrapy/utils/_stopmode.py diff --git a/scrapy/core/downloader/__init__.py b/scrapy/core/downloader/__init__.py index 905b32d08..9b4b11f87 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,9 @@ class Downloader: DownloaderMiddlewareManager.from_crawler(crawler) ) 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" ) @@ -174,6 +177,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) @@ -208,12 +215,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) @@ -255,11 +271,49 @@ 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 + self._fast_stopping = True + + 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 94920e840..f577b481b 100644 --- a/scrapy/core/engine.py +++ b/scrapy/core/engine.py @@ -24,10 +24,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 StopMode, max_stop_mode, normalize_stop_mode from scrapy.utils.asyncio import ( AsyncioLoopingCall, create_looping_call, @@ -118,6 +120,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 @@ -201,23 +205,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: @@ -232,7 +249,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) @@ -403,6 +420,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) @@ -582,20 +610,50 @@ 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}, + ) + + 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") @@ -603,6 +661,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 @@ -615,6 +675,9 @@ class ExecutionEngine: def log_failure(msg: str) -> None: logger.error(msg, exc_info=True, extra={"spider": spider}) # noqa: LOG014 + 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 0a19e9985..2915e375b 100644 --- a/scrapy/crawler.py +++ b/scrapy/crawler.py @@ -19,7 +19,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 StopMode, normalize_stop_mode +from scrapy.utils.defer import deferred_from_coro, ensure_awaitable from scrapy.utils.log import ( configure_logging, get_scrapy_root_handler, @@ -41,7 +42,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 @@ -84,6 +85,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: @@ -232,7 +242,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( @@ -240,18 +250,39 @@ 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 + + # Keep repeated graceful stop calls as no-ops, while allowing + # fast/force escalation when shutdown is already in progress. + 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 + + try: + await self.engine.stop_async(mode=mode) + except RuntimeError as exc: + if str(exc) != "Engine not running": + raise @staticmethod def _get_component( @@ -470,13 +501,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]: @@ -592,15 +626,16 @@ class AsyncCrawlerRunner(CrawlerRunnerBase): task.add_done_callback(_done) 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: @@ -622,6 +657,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 @@ -631,10 +674,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 @@ -646,7 +696,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}, ) @@ -654,7 +712,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: @@ -696,13 +755,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() @@ -755,10 +821,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 @@ -855,8 +923,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 @@ -1001,14 +1079,12 @@ 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 def _create_shutdown_task() -> None: - coro = self._shutdown_graceful_reactorless() + coro = self._shutdown_reactorless(mode=mode) try: loop.create_task(coro) except RuntimeError: @@ -1016,8 +1092,18 @@ class AsyncCrawlerProcess(CrawlerProcessBase, AsyncCrawlerRunner): 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..afb7b4768 --- /dev/null +++ b/scrapy/utils/_stopmode.py @@ -0,0 +1,30 @@ +from __future__ import annotations + +from typing import Literal, cast + +StopMode = Literal["graceful", "fast", "force"] + +_STOP_MODE_PRIORITY: dict[StopMode, int] = { + "graceful": 0, + "fast": 1, + "force": 2, +} + + +def normalize_stop_mode(mode: StopMode | None, *, allow_force: bool = True) -> StopMode: + if mode is None: + return "graceful" + if mode not in _STOP_MODE_PRIORITY: + raise ValueError( + f"Unknown stop mode {mode!r}. Expected one of: graceful, fast, force" + ) + normalized = cast("StopMode", mode) + if normalized == "force" and not allow_force: + raise ValueError("The force stop mode is not supported in this context") + return normalized + + +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/test_core_downloader.py b/tests/test_core_downloader.py index f9c5abf8a..c3ba4125c 100644 --- a/tests/test_core_downloader.py +++ b/tests/test_core_downloader.py @@ -6,17 +6,19 @@ from typing import TYPE_CHECKING, cast import OpenSSL.SSL import pytest from pytest_twisted import async_yield_fixture +from twisted.internet.defer import Deferred 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 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 from scrapy.utils.defer import maybe_deferred_to_future from scrapy.utils.misc import build_from_crawler @@ -28,7 +30,6 @@ 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.internet.ssl import ContextFactory from twisted.web.iweb import IBodyProducer @@ -207,3 +208,38 @@ async def test_fetch_deprecated_spider_arg(): match=r"The fetch\(\) method of .+\.CustomDownloader requires a spider argument", ): 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 = Deferred() + failures = [] + 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")) diff --git a/tests/test_crawler.py b/tests/test_crawler.py index cbcb7e274..6ae097664 100644 --- a/tests/test_crawler.py +++ b/tests/test_crawler.py @@ -788,3 +788,45 @@ 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") + + +@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 + + crawler.engine = DummyEngine() # type: ignore[assignment] + + with caplog.at_level(logging.WARNING): + await crawler.stop_async(mode="force") + + assert crawler.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 diff --git a/tests/test_crawler_subprocess.py b/tests/test_crawler_subprocess.py index beae4f277..ec40d0488 100644 --- a/tests/test_crawler_subprocess.py +++ b/tests/test_crawler_subprocess.py @@ -227,7 +227,10 @@ 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("dropping downloader requests") await async_sleep(0.01) p.kill(sig) p.expect_exact("forcing unclean shutdown") diff --git a/tests/test_engine.py b/tests/test_engine.py index 2cd583721..78e52d7db 100644 --- a/tests/test_engine.py +++ b/tests/test_engine.py @@ -768,3 +768,26 @@ 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 + + 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 From 15af004382ccd97c87455412f9d05ee10620bccf Mon Sep 17 00:00:00 2001 From: Adrian Chaves Date: Mon, 27 Apr 2026 15:28:27 +0200 Subject: [PATCH 02/18] Address typing and linting issues --- scrapy/core/engine.py | 1 + scrapy/utils/_stopmode.py | 7 +++---- tests/test_core_downloader.py | 3 ++- tests/test_crawler.py | 7 ++++--- tests/test_engine.py | 7 +++---- 5 files changed, 13 insertions(+), 12 deletions(-) diff --git a/scrapy/core/engine.py b/scrapy/core/engine.py index f577b481b..973ae3258 100644 --- a/scrapy/core/engine.py +++ b/scrapy/core/engine.py @@ -641,6 +641,7 @@ class ExecutionEngine: extra={"spider": self.spider}, ) + # pylint: disable=too-many-statements async def close_spider_async( # noqa: PLR0912, PLR0915 self, *, diff --git a/scrapy/utils/_stopmode.py b/scrapy/utils/_stopmode.py index afb7b4768..05b64ec9a 100644 --- a/scrapy/utils/_stopmode.py +++ b/scrapy/utils/_stopmode.py @@ -1,6 +1,6 @@ from __future__ import annotations -from typing import Literal, cast +from typing import Literal StopMode = Literal["graceful", "fast", "force"] @@ -18,10 +18,9 @@ def normalize_stop_mode(mode: StopMode | None, *, allow_force: bool = True) -> S raise ValueError( f"Unknown stop mode {mode!r}. Expected one of: graceful, fast, force" ) - normalized = cast("StopMode", mode) - if normalized == "force" and not allow_force: + if mode == "force" and not allow_force: raise ValueError("The force stop mode is not supported in this context") - return normalized + return mode def max_stop_mode(mode1: StopMode, mode2: StopMode) -> StopMode: diff --git a/tests/test_core_downloader.py b/tests/test_core_downloader.py index c3ba4125c..aba1541e0 100644 --- a/tests/test_core_downloader.py +++ b/tests/test_core_downloader.py @@ -31,6 +31,7 @@ from tests.utils.decorators import coroutine_test if TYPE_CHECKING: from twisted.internet.ssl import ContextFactory + from twisted.python.failure import Failure from twisted.web.iweb import IBodyProducer @@ -220,7 +221,7 @@ async def test_stop_async_drops_queued_requests() -> None: request = Request("https://example.com") queue_dfd: Deferred = Deferred() - failures = [] + failures: list[Failure] = [] queue_dfd.addErrback(failures.append) slot.queue.append((request, queue_dfd)) diff --git a/tests/test_crawler.py b/tests/test_crawler.py index 6ae097664..c4c28114c 100644 --- a/tests/test_crawler.py +++ b/tests/test_crawler.py @@ -794,7 +794,7 @@ async def test_deprecated_crawler_stop() -> None: 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") + await crawler.stop_async(mode="invalid") # type: ignore[arg-type] @coroutine_test @@ -809,12 +809,13 @@ async def test_crawler_force_stop_falls_back_to_fast( async def stop_async(self, *, mode: str = "graceful") -> None: self.called_mode = mode - crawler.engine = DummyEngine() # type: ignore[assignment] + dummy_engine = DummyEngine() + crawler.engine = dummy_engine # type: ignore[assignment] with caplog.at_level(logging.WARNING): await crawler.stop_async(mode="force") - assert crawler.engine.called_mode == "fast" + assert dummy_engine.called_mode == "fast" assert "Falling back to fast stop" in caplog.text diff --git a/tests/test_engine.py b/tests/test_engine.py index 78e52d7db..5900695af 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 +from unittest.mock import Mock, call, patch from urllib.parse import urlparse import attr @@ -784,9 +784,8 @@ class TestEngineCloseSpider: calls += 1 return 3 - engine.downloader.stop_async = fast_stop_downloader - - await engine.close_spider_async(mode="fast") + with patch.object(engine.downloader, "stop_async", fast_stop_downloader): + await engine.close_spider_async(mode="fast") assert calls == 1 assert crawler.stats From 44f0fd0d37832f2af1b5ddeb246b7e2490dce82d Mon Sep 17 00:00:00 2001 From: Adrian Chaves Date: Wed, 29 Apr 2026 12:25:54 +0200 Subject: [PATCH 03/18] Solve issue making CI hang? --- scrapy/crawler.py | 5 +++++ tests/test_crawler.py | 20 ++++++++++++++++++++ 2 files changed, 25 insertions(+) diff --git a/scrapy/crawler.py b/scrapy/crawler.py index 2915e375b..ace5f4237 100644 --- a/scrapy/crawler.py +++ b/scrapy/crawler.py @@ -278,6 +278,11 @@ class Crawler: 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: diff --git a/tests/test_crawler.py b/tests/test_crawler.py index c4c28114c..75a02aadc 100644 --- a/tests/test_crawler.py +++ b/tests/test_crawler.py @@ -797,6 +797,26 @@ async def test_crawler_stop_async_invalid_mode() -> None: 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, From 03f89c99ab355eca0e94a862c78ad65c94238e70 Mon Sep 17 00:00:00 2001 From: Adrian Chaves Date: Wed, 29 Apr 2026 12:45:16 +0200 Subject: [PATCH 04/18] Remove unused code and privatize to protect against import of imports from public modules --- scrapy/core/engine.py | 16 ++++++++-------- scrapy/crawler.py | 26 +++++++++++++------------- scrapy/utils/_stopmode.py | 10 ++++------ 3 files changed, 25 insertions(+), 27 deletions(-) diff --git a/scrapy/core/engine.py b/scrapy/core/engine.py index 973ae3258..d9e0c5eed 100644 --- a/scrapy/core/engine.py +++ b/scrapy/core/engine.py @@ -29,7 +29,7 @@ from scrapy.exceptions import ( ScrapyDeprecationWarning, ) from scrapy.http import Request, Response -from scrapy.utils._stopmode import StopMode, max_stop_mode, normalize_stop_mode +from scrapy.utils._stopmode import _normalize_stop_mode, _StopMode, max_stop_mode from scrapy.utils.asyncio import ( AsyncioLoopingCall, create_looping_call, @@ -120,7 +120,7 @@ class ExecutionEngine: self.running: bool = False self._starting: bool = False self._stopping: bool = False - self._stop_mode: StopMode = "graceful" + self._stop_mode: _StopMode = "graceful" self._downloader_fast_stopped: bool = False self.paused: bool = False self._spider_closed_callback: Callable[ @@ -206,7 +206,7 @@ class ExecutionEngine: await maybe_deferred_to_future(self._closewait) def stop( - self, *, mode: StopMode = "graceful" + self, *, mode: _StopMode = "graceful" ) -> Deferred[None]: # pragma: no cover warnings.warn( "ExecutionEngine.stop() is deprecated, use stop_async() instead", @@ -215,13 +215,13 @@ class ExecutionEngine: ) return deferred_from_coro(self.stop_async(mode=mode)) - async def stop_async(self, *, mode: StopMode = "graceful") -> None: + async def stop_async(self, *, mode: _StopMode = "graceful") -> None: """Gracefully stop the execution engine. .. versionadded:: 2.14 """ - mode = normalize_stop_mode(mode, allow_force=False) + mode = _normalize_stop_mode(mode, allow_force=False) if not self._starting and not self._stopping: raise RuntimeError("Engine not running") @@ -613,7 +613,7 @@ class ExecutionEngine: self, spider: Spider, reason: str = "cancelled", - mode: StopMode = "graceful", + mode: _StopMode = "graceful", ) -> Deferred[None]: # pragma: no cover warnings.warn( "ExecutionEngine.close_spider() is deprecated, use close_spider_async() instead", @@ -646,13 +646,13 @@ class ExecutionEngine: self, *, reason: str = "cancelled", - mode: StopMode = "graceful", + mode: _StopMode = "graceful", ) -> None: """Close (cancel) spider and clear all its outstanding requests. .. versionadded:: 2.14 """ - mode = normalize_stop_mode(mode, allow_force=False) + mode = _normalize_stop_mode(mode, allow_force=False) self._stop_mode = max_stop_mode(self._stop_mode, mode) if self.spider is None: diff --git a/scrapy/crawler.py b/scrapy/crawler.py index ace5f4237..7fd39571e 100644 --- a/scrapy/crawler.py +++ b/scrapy/crawler.py @@ -19,7 +19,7 @@ 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._stopmode import StopMode, normalize_stop_mode +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, @@ -242,7 +242,7 @@ class Crawler: def _create_engine(self) -> ExecutionEngine: return ExecutionEngine(self, lambda _: self.stop_async()) - def stop(self, *, mode: StopMode = "graceful") -> 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( @@ -252,12 +252,12 @@ class Crawler: ) return deferred_from_coro(self.stop_async(mode=mode)) - async def stop_async(self, *, mode: StopMode = "graceful") -> 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 """ - mode = normalize_stop_mode(mode) + mode = _normalize_stop_mode(mode) was_crawling = self.crawling self.crawling = False @@ -506,13 +506,13 @@ class CrawlerRunner(CrawlerRunnerBase): self._active.discard(d) self.bootstrap_failed |= not getattr(crawler, "spider", None) or failed - def stop(self, *, mode: StopMode = "graceful") -> 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. """ - mode = normalize_stop_mode(mode) + mode = _normalize_stop_mode(mode) return DeferredList( deferred_from_coro(c.stop_async(mode=mode)) for c in self.crawlers ) @@ -631,13 +631,13 @@ class AsyncCrawlerRunner(CrawlerRunnerBase): task.add_done_callback(_done) return task - async def stop(self, *, mode: StopMode = "graceful") -> 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) + mode = _normalize_stop_mode(mode) if self.crawlers: await asyncio.wait( [asyncio.create_task(c.stop_async(mode=mode)) for c in self.crawlers] @@ -760,7 +760,7 @@ class CrawlerProcessBase(CrawlerRunnerBase): ) @abstractmethod - def _stop_dfd(self, *, mode: StopMode = "graceful") -> Deferred[Any]: + def _stop_dfd(self, *, mode: _StopMode = "graceful") -> Deferred[Any]: raise NotImplementedError @inlineCallbacks @@ -830,7 +830,7 @@ class CrawlerProcess(CrawlerProcessBase, CrawlerRunner): crawler._set_force_stop_callback(self._force_stop) return crawler - def _stop_dfd(self, *, mode: StopMode = "graceful") -> Deferred[Any]: + def _stop_dfd(self, *, mode: _StopMode = "graceful") -> Deferred[Any]: return self.stop(mode=mode) def start( @@ -938,7 +938,7 @@ class AsyncCrawlerProcess(CrawlerProcessBase, AsyncCrawlerRunner): if (task := self._reactorless_main_task) is not None: loop.call_soon_threadsafe(task.cancel) - def _stop_dfd(self, *, mode: StopMode = "graceful") -> Deferred[Any]: + def _stop_dfd(self, *, mode: _StopMode = "graceful") -> Deferred[Any]: return deferred_from_coro(self.stop(mode=mode)) def start( @@ -1084,7 +1084,7 @@ class AsyncCrawlerProcess(CrawlerProcessBase, AsyncCrawlerRunner): } ) - def _schedule_reactorless_shutdown(self, *, mode: StopMode) -> None: + def _schedule_reactorless_shutdown(self, *, mode: _StopMode) -> None: if (loop := self._reactorless_loop) is None: return @@ -1107,7 +1107,7 @@ class AsyncCrawlerProcess(CrawlerProcessBase, AsyncCrawlerRunner): self._log_fast_shutdown(signum) self._schedule_reactorless_shutdown(mode="fast") - async def _shutdown_reactorless(self, *, mode: StopMode) -> None: + 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 diff --git a/scrapy/utils/_stopmode.py b/scrapy/utils/_stopmode.py index 05b64ec9a..6acd1c42e 100644 --- a/scrapy/utils/_stopmode.py +++ b/scrapy/utils/_stopmode.py @@ -2,18 +2,16 @@ from __future__ import annotations from typing import Literal -StopMode = Literal["graceful", "fast", "force"] +_StopMode = Literal["graceful", "fast", "force"] -_STOP_MODE_PRIORITY: dict[StopMode, int] = { +_STOP_MODE_PRIORITY: dict[_StopMode, int] = { "graceful": 0, "fast": 1, "force": 2, } -def normalize_stop_mode(mode: StopMode | None, *, allow_force: bool = True) -> StopMode: - if mode is None: - return "graceful" +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" @@ -23,7 +21,7 @@ def normalize_stop_mode(mode: StopMode | None, *, allow_force: bool = True) -> S return mode -def max_stop_mode(mode1: StopMode, mode2: StopMode) -> StopMode: +def max_stop_mode(mode1: _StopMode, mode2: _StopMode) -> _StopMode: if _STOP_MODE_PRIORITY[mode1] >= _STOP_MODE_PRIORITY[mode2]: return mode1 return mode2 From 41fa1e94a07286fd596e2b45fa4f2ba26f809402 Mon Sep 17 00:00:00 2001 From: Adrian Chaves Date: Wed, 29 Apr 2026 12:45:52 +0200 Subject: [PATCH 05/18] Add CI timeouts while at it --- .github/workflows/checks.yml | 2 ++ .github/workflows/tests-macos.yml | 1 + .github/workflows/tests-ubuntu.yml | 1 + .github/workflows/tests-windows.yml | 1 + 4 files changed, 5 insertions(+) diff --git a/.github/workflows/checks.yml b/.github/workflows/checks.yml index cb05784fa..f33844c41 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 2a1c62833..9f60310a5 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 a193ca05f..3fcca5404 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 48aa56e15..cf3102d31 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: From 17c406d91b6606b9fe9436e434ea61e75dc7a0e1 Mon Sep 17 00:00:00 2001 From: Adrian Chaves Date: Wed, 29 Apr 2026 13:12:41 +0200 Subject: [PATCH 06/18] Improve coverage --- scrapy/core/downloader/__init__.py | 2 - tests/test_core_downloader.py | 81 +++++++++++++++++++++++++++++- tests/test_crawler.py | 31 ++++++++++++ tests/test_engine.py | 78 ++++++++++++++++++++++++++-- 4 files changed, 185 insertions(+), 7 deletions(-) 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() From fe8f527b1a5b91146351a4752d62265f1c3a55a7 Mon Sep 17 00:00:00 2001 From: Adrian Chaves Date: Wed, 29 Apr 2026 13:39:39 +0200 Subject: [PATCH 07/18] Fix test for the default reactor --- tests/test_engine.py | 9 +++++---- 1 file changed, 5 insertions(+), 4 deletions(-) diff --git a/tests/test_engine.py b/tests/test_engine.py index 4c408d13c..54366fba7 100644 --- a/tests/test_engine.py +++ b/tests/test_engine.py @@ -37,6 +37,7 @@ 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: @@ -474,14 +475,14 @@ class TestEngine(TestEngineBase): 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) + 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 task.done() + assert not stop_dfd.called assert engine._closewait engine._closewait.callback(None) - await task + await maybe_deferred_to_future(stop_dfd) @coroutine_test async def test_handle_downloader_output_ignores_fast_cancelled_failures( From dbcd87c09c03e01883bfff5efebb809c581c1431 Mon Sep 17 00:00:00 2001 From: Adrian Chaves Date: Wed, 29 Apr 2026 14:52:17 +0200 Subject: [PATCH 08/18] Improve test coverage --- tests/test_crawler.py | 114 ++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 114 insertions(+) diff --git a/tests/test_crawler.py b/tests/test_crawler.py index 9cc9a9d5d..298445f84 100644 --- a/tests/test_crawler.py +++ b/tests/test_crawler.py @@ -882,3 +882,117 @@ async def test_crawler_stop_async_ignores_engine_not_running_runtime_error() -> 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() -> 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() -> 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_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 = object() + 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 = loop + crawler_process._reactorless_main_task = task + crawler = crawler_process.create_crawler(DefaultSpider) + + await crawler.stop_async(mode="force") + + assert loop.callback == task.cancel From 2a0825d75a8873a3d2bacb359986ff334305b417 Mon Sep 17 00:00:00 2001 From: Adrian Chaves Date: Wed, 29 Apr 2026 15:18:05 +0200 Subject: [PATCH 09/18] Address typing and linting issues --- tests/test_crawler.py | 14 ++++++++------ 1 file changed, 8 insertions(+), 6 deletions(-) diff --git a/tests/test_crawler.py b/tests/test_crawler.py index 298445f84..34f5aa18a 100644 --- a/tests/test_crawler.py +++ b/tests/test_crawler.py @@ -907,7 +907,7 @@ async def test_crawler_process_force_stop_via_public_crawler_api() -> None: crawler_process = CrawlerProcess(install_root_handler=False) called = False - def stop_reactor() -> None: + def stop_reactor(_: Any = None) -> None: nonlocal called called = True @@ -931,7 +931,7 @@ async def test_async_crawler_process_force_stop_reactor_enabled_via_public_crawl ) called = False - def stop_reactor() -> None: + def stop_reactor(_: Any = None) -> None: nonlocal called called = True @@ -956,7 +956,7 @@ async def test_async_crawler_process_force_stop_reactorless_without_loop( install_root_handler=False, ) crawler_process._reactorless_loop = None - crawler_process._reactorless_main_task = object() + crawler_process._reactorless_main_task = None crawler = crawler_process.create_crawler(DefaultSpider) await crawler.stop_async(mode="force") @@ -989,10 +989,12 @@ async def test_async_crawler_process_force_stop_reactorless_with_task( loop = DummyLoop() task = DummyTask() - crawler_process._reactorless_loop = loop - crawler_process._reactorless_main_task = task + 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 == task.cancel + assert loop.callback is not None + loop.callback() + assert task.called is True From de950c9fa10feae962f575478ec3308403da5da7 Mon Sep 17 00:00:00 2001 From: Adrian Chaves Date: Wed, 29 Apr 2026 16:06:59 +0200 Subject: [PATCH 10/18] Improve coverage --- ...eactorless_sleeping_no_stop_after_crawl.py | 20 ++++++ tests/test_crawler.py | 64 +++++++++++++++++++ tests/test_crawler_subprocess.py | 3 + tests/test_engine.py | 14 ++++ 4 files changed, 101 insertions(+) create mode 100644 tests/AsyncCrawlerProcess/reactorless_sleeping_no_stop_after_crawl.py 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_crawler.py b/tests/test_crawler.py index 34f5aa18a..d0ead784b 100644 --- a/tests/test_crawler.py +++ b/tests/test_crawler.py @@ -943,6 +943,25 @@ async def test_async_crawler_process_force_stop_reactor_enabled_via_public_crawl 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( @@ -998,3 +1017,48 @@ async def test_async_crawler_process_force_stop_reactorless_with_task( 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() + + def shutdown_reactorless(*, mode: str) -> DummyCoro: + assert mode == "graceful" + return coro + + crawler_process._reactorless_loop = cast("asyncio.AbstractEventLoop", loop) + crawler_process._shutdown_reactorless = shutdown_reactorless + + crawler_process._schedule_reactorless_shutdown(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 ec40d0488..4893537f8 100644 --- a/tests/test_crawler_subprocess.py +++ b/tests/test_crawler_subprocess.py @@ -420,6 +420,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 54366fba7..38047a35c 100644 --- a/tests/test_engine.py +++ b/tests/test_engine.py @@ -484,6 +484,20 @@ class TestEngine(TestEngineBase): 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, From dc4db5318f9ee6c1f9460110c60d6ab28b015188 Mon Sep 17 00:00:00 2001 From: Adrian Chaves Date: Wed, 29 Apr 2026 16:23:23 +0200 Subject: [PATCH 11/18] Address typing issue --- tests/test_crawler.py | 12 +++++++++--- 1 file changed, 9 insertions(+), 3 deletions(-) diff --git a/tests/test_crawler.py b/tests/test_crawler.py index d0ead784b..c01a37bd9 100644 --- a/tests/test_crawler.py +++ b/tests/test_crawler.py @@ -5,6 +5,7 @@ import warnings from collections.abc import Generator from pathlib import Path from typing import Any, cast +from unittest.mock import patch import pytest from twisted.internet.defer import Deferred @@ -1049,15 +1050,20 @@ def test_async_crawler_process_schedule_reactorless_shutdown_runtime_error() -> loop = DummyLoop() coro = DummyCoro() + called_mode: str | None = None def shutdown_reactorless(*, mode: str) -> DummyCoro: - assert mode == "graceful" + nonlocal called_mode + called_mode = mode return coro crawler_process._reactorless_loop = cast("asyncio.AbstractEventLoop", loop) - crawler_process._shutdown_reactorless = shutdown_reactorless + with patch.object( + crawler_process, "_shutdown_reactorless", new=shutdown_reactorless + ): + crawler_process._schedule_reactorless_shutdown(mode="graceful") - crawler_process._schedule_reactorless_shutdown(mode="graceful") + assert called_mode == "graceful" assert loop.scheduled assert loop.create_task_called From f5080c1e4c8ece4f51745f9cae6f49b640dacea3 Mon Sep 17 00:00:00 2001 From: Adrian Chaves Date: Wed, 29 Apr 2026 16:38:32 +0200 Subject: [PATCH 12/18] Minor refactoring --- scrapy/core/engine.py | 6 +++--- scrapy/utils/_stopmode.py | 2 +- 2 files changed, 4 insertions(+), 4 deletions(-) diff --git a/scrapy/core/engine.py b/scrapy/core/engine.py index d9e0c5eed..9416783f1 100644 --- a/scrapy/core/engine.py +++ b/scrapy/core/engine.py @@ -29,7 +29,7 @@ from scrapy.exceptions import ( ScrapyDeprecationWarning, ) from scrapy.http import Request, Response -from scrapy.utils._stopmode import _normalize_stop_mode, _StopMode, max_stop_mode +from scrapy.utils._stopmode import _max_stop_mode, _normalize_stop_mode, _StopMode from scrapy.utils.asyncio import ( AsyncioLoopingCall, create_looping_call, @@ -226,7 +226,7 @@ class ExecutionEngine: if not self._starting and not self._stopping: raise RuntimeError("Engine not running") - self._stop_mode = max_stop_mode(self._stop_mode, mode) + self._stop_mode = _max_stop_mode(self._stop_mode, mode) if self._stopping: if self.spider is not None and self._stop_mode == "fast": @@ -653,7 +653,7 @@ class ExecutionEngine: .. versionadded:: 2.14 """ mode = _normalize_stop_mode(mode, allow_force=False) - self._stop_mode = max_stop_mode(self._stop_mode, mode) + self._stop_mode = _max_stop_mode(self._stop_mode, mode) if self.spider is None: raise RuntimeError("Spider not opened") diff --git a/scrapy/utils/_stopmode.py b/scrapy/utils/_stopmode.py index 6acd1c42e..ec4d083a6 100644 --- a/scrapy/utils/_stopmode.py +++ b/scrapy/utils/_stopmode.py @@ -21,7 +21,7 @@ def _normalize_stop_mode(mode: _StopMode, *, allow_force: bool = True) -> _StopM return mode -def max_stop_mode(mode1: _StopMode, mode2: _StopMode) -> _StopMode: +def _max_stop_mode(mode1: _StopMode, mode2: _StopMode) -> _StopMode: if _STOP_MODE_PRIORITY[mode1] >= _STOP_MODE_PRIORITY[mode2]: return mode1 return mode2 From 1b917950f55024d062a7675602ef904c34a1890c Mon Sep 17 00:00:00 2001 From: Adrian Chaves Date: Wed, 29 Apr 2026 16:40:41 +0200 Subject: [PATCH 13/18] Remove unneeded comment --- scrapy/crawler.py | 2 -- 1 file changed, 2 deletions(-) diff --git a/scrapy/crawler.py b/scrapy/crawler.py index 7fd39571e..671d49d9a 100644 --- a/scrapy/crawler.py +++ b/scrapy/crawler.py @@ -261,8 +261,6 @@ class Crawler: was_crawling = self.crawling self.crawling = False - # Keep repeated graceful stop calls as no-ops, while allowing - # fast/force escalation when shutdown is already in progress. if not was_crawling and mode == "graceful": return From c5250f1d803d2a1bb516dc36f3e816cc927db4ea Mon Sep 17 00:00:00 2001 From: Adrian Chaves Date: Wed, 29 Apr 2026 21:43:48 +0200 Subject: [PATCH 14/18] Remove sleep that may have been making a test brittle --- tests/test_crawler_subprocess.py | 1 - 1 file changed, 1 deletion(-) diff --git a/tests/test_crawler_subprocess.py b/tests/test_crawler_subprocess.py index 4893537f8..b12132a82 100644 --- a/tests/test_crawler_subprocess.py +++ b/tests/test_crawler_subprocess.py @@ -231,7 +231,6 @@ class TestCrawlerProcessSubprocessBase(ScriptRunnerMixin): await async_sleep(0.01) p.kill(sig) p.expect_exact("dropping downloader requests") - await async_sleep(0.01) p.kill(sig) p.expect_exact("forcing unclean shutdown") p.wait() # type: ignore[no-untyped-call] From 1e6aca0bff900b561e55f74deb44a41db7740bbf Mon Sep 17 00:00:00 2001 From: Adrian Chaves Date: Wed, 29 Apr 2026 22:22:29 +0200 Subject: [PATCH 15/18] =?UTF-8?q?=F0=9F=A4=9E?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- tests/test_crawler_subprocess.py | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/tests/test_crawler_subprocess.py b/tests/test_crawler_subprocess.py index b12132a82..a5f1917e2 100644 --- a/tests/test_crawler_subprocess.py +++ b/tests/test_crawler_subprocess.py @@ -11,6 +11,7 @@ from typing import TYPE_CHECKING import pytest from packaging.version import parse as parse_version +from pexpect.exceptions import EOF from pexpect.popen_spawn import PopenSpawn from w3lib import __version__ as w3lib_version @@ -232,7 +233,9 @@ class TestCrawlerProcessSubprocessBase(ScriptRunnerMixin): p.kill(sig) p.expect_exact("dropping downloader requests") p.kill(sig) - p.expect_exact("forcing unclean shutdown") + # Depending on timing, fast shutdown may complete before the third + # signal handler logs the force-shutdown message. + p.expect_exact(["forcing unclean shutdown", EOF]) p.wait() # type: ignore[no-untyped-call] @coroutine_test From c1e3997ec0c35dcf9b64e34bf6b41f0554cfa6e7 Mon Sep 17 00:00:00 2001 From: Adrian Chaves Date: Wed, 29 Apr 2026 22:55:03 +0200 Subject: [PATCH 16/18] Try simply increasing the timeout --- tests/test_crawler_subprocess.py | 6 ++---- 1 file changed, 2 insertions(+), 4 deletions(-) diff --git a/tests/test_crawler_subprocess.py b/tests/test_crawler_subprocess.py index a5f1917e2..b3fa101d9 100644 --- a/tests/test_crawler_subprocess.py +++ b/tests/test_crawler_subprocess.py @@ -11,7 +11,6 @@ from typing import TYPE_CHECKING import pytest from packaging.version import parse as parse_version -from pexpect.exceptions import EOF from pexpect.popen_spawn import PopenSpawn from w3lib import __version__ as w3lib_version @@ -232,10 +231,9 @@ class TestCrawlerProcessSubprocessBase(ScriptRunnerMixin): await async_sleep(0.01) p.kill(sig) p.expect_exact("dropping downloader requests") + await async_sleep(0.01) p.kill(sig) - # Depending on timing, fast shutdown may complete before the third - # signal handler logs the force-shutdown message. - p.expect_exact(["forcing unclean shutdown", EOF]) + p.expect_exact("forcing unclean shutdown", timeout=10) p.wait() # type: ignore[no-untyped-call] @coroutine_test From d8ac397be1e3eb4ae8f900ce1e098ac431b94504 Mon Sep 17 00:00:00 2001 From: Adrian Chaves Date: Wed, 29 Apr 2026 23:07:44 +0200 Subject: [PATCH 17/18] Number big not boom? --- tests/test_crawler_subprocess.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/tests/test_crawler_subprocess.py b/tests/test_crawler_subprocess.py index b3fa101d9..ef7b3dd0a 100644 --- a/tests/test_crawler_subprocess.py +++ b/tests/test_crawler_subprocess.py @@ -233,7 +233,7 @@ class TestCrawlerProcessSubprocessBase(ScriptRunnerMixin): p.expect_exact("dropping downloader requests") await async_sleep(0.01) p.kill(sig) - p.expect_exact("forcing unclean shutdown", timeout=10) + p.expect_exact("forcing unclean shutdown", timeout=20) p.wait() # type: ignore[no-untyped-call] @coroutine_test From 821afd58044716375d71822db4d7a02e4ae22a07 Mon Sep 17 00:00:00 2001 From: Adrian Chaves Date: Tue, 16 Jun 2026 17:21:06 +0200 Subject: [PATCH 18/18] Address typing issues --- tests/test_core_downloader.py | 14 +++++++------- 1 file changed, 7 insertions(+), 7 deletions(-) diff --git a/tests/test_core_downloader.py b/tests/test_core_downloader.py index 8877f2874..c16725356 100644 --- a/tests/test_core_downloader.py +++ b/tests/test_core_downloader.py @@ -1,7 +1,7 @@ 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 @@ -279,7 +279,7 @@ async def test_stop_async_drops_queued_requests() -> None: downloader.slots["example.com"] = slot request = Request("https://example.com") - queue_dfd: Deferred = Deferred() + queue_dfd: Deferred[Any] = Deferred() failures: list[Failure] = [] queue_dfd.addErrback(failures.append) slot.queue.append((request, queue_dfd)) @@ -311,7 +311,7 @@ async def test_wait_for_download_errbacks_queue_deferred_on_error() -> None: downloader = Downloader(crawler) slot = Slot(concurrency=1, delay=0, randomize_delay=False) - queue_dfd: Deferred = Deferred() + queue_dfd: Deferred[Any] = Deferred() failures: list[Failure] = [] queue_dfd.addErrback(failures.append) @@ -332,7 +332,7 @@ async def test_wait_for_download_keeps_called_queue_deferred_on_error() -> None: downloader = Downloader(crawler) slot = Slot(concurrency=1, delay=0, randomize_delay=False) - queue_dfd: Deferred = Deferred() + queue_dfd: Deferred[Any] = Deferred() queue_dfd.callback(None) with patch.object(downloader, "_download", side_effect=RuntimeError("boom")): @@ -353,7 +353,7 @@ async def test_stop_async_skips_called_queued_deferred() -> None: slot = Slot(concurrency=1, delay=0, randomize_delay=False) downloader.slots["example.com"] = slot - queue_dfd: Deferred = Deferred() + queue_dfd: Deferred[Any] = Deferred() queue_dfd.callback(None) slot.queue.append((Request("https://example.com"), queue_dfd)) @@ -366,10 +366,10 @@ async def test_stop_async_cancels_pending_download_tasks() -> None: crawler = get_crawler(DefaultSpider) downloader = Downloader(crawler) - done_dfd: Deferred = Deferred() + done_dfd: Deferred[None] = Deferred() done_dfd.callback(None) - pending_dfd: Deferred = Deferred() + pending_dfd: Deferred[None] = Deferred() failures: list[Failure] = [] pending_dfd.addErrback(failures.append)