This commit is contained in:
Adrian 2026-08-15 11:46:48 -05:00 committed by GitHub
commit a58f4be0d3
No known key found for this signature in database
GPG Key ID: B5690EEEBB952194
14 changed files with 837 additions and 74 deletions

View File

@ -17,6 +17,7 @@ concurrency:
jobs:
checks:
runs-on: ubuntu-latest
timeout-minutes: 30
env:
# Make uv use the interpreter that actions/setup-python installed instead
# of downloading one of its own.
@ -69,6 +70,7 @@ jobs:
pre-commit:
runs-on: ubuntu-latest
timeout-minutes: 15
steps:
- uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1
with:

View File

@ -18,6 +18,7 @@ jobs:
tests:
name: tests (${{ matrix.python-version }}, ${{ matrix.env.TOXENV }})
runs-on: macos-latest
timeout-minutes: 60
env:
PYTEST_ADDOPTS: ${{ matrix.coverage && '-n auto' || '-n auto --no-cov' }}
# Make uv use the interpreter that actions/setup-python installed instead

View File

@ -18,6 +18,7 @@ jobs:
tests:
name: tests (${{ matrix.python-version }}, ${{ matrix.env.TOXENV }})
runs-on: ubuntu-latest
timeout-minutes: 60
env:
PYTEST_ADDOPTS: ${{ matrix.coverage && '-n auto' || '-n auto --no-cov' }}
# Make uv use the interpreter that actions/setup-python installed instead

View File

@ -18,6 +18,7 @@ jobs:
tests:
name: tests (${{ matrix.python-version }}, ${{ matrix.env.TOXENV }})
runs-on: windows-latest
timeout-minutes: 60
env:
PYTEST_ADDOPTS: ${{ matrix.coverage && '-n auto' || '-n auto --no-cov' }}
# Make uv use the interpreter that actions/setup-python installed instead

View File

@ -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,
)
@ -104,6 +104,8 @@ class Downloader:
DownloaderMiddlewareManager, 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"
)
@ -159,6 +161,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)
@ -193,12 +199,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)
@ -240,11 +255,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()

View File

@ -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,
@ -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
@ -204,23 +208,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:
@ -235,7 +252,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)
@ -420,6 +437,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)
@ -616,20 +644,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")
@ -637,6 +696,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
@ -646,6 +707,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:

View File

@ -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,
@ -46,7 +47,7 @@ from scrapy.utils.reactorless import (
)
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
@ -132,6 +133,9 @@ class Crawler:
self._started: bool = False
self.spider: Spider | None = None
self._force_stop_callback: (
Callable[[], Awaitable[None] | Deferred[None] | None] | None
) = None
self._engine: ExecutionEngine | None = None
self._extensions: ExtensionManager | None = None
@ -139,6 +143,12 @@ class Crawler:
self._request_fingerprinter: RequestFingerprinterProtocol | None = None
self._stats: StatsCollector | 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:
# scrapy root handler already installed: update it with new settings
@ -322,7 +332,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(
@ -330,17 +340,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
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(
@ -563,13 +598,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]:
@ -687,15 +725,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:
@ -717,6 +756,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
@ -726,10 +773,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
@ -741,7 +795,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},
)
@ -749,7 +811,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:
@ -791,13 +854,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()
@ -850,10 +920,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
@ -951,8 +1023,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
@ -1101,24 +1183,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()

27
scrapy/utils/_stopmode.py Normal file
View File

@ -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

View File

@ -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)

View File

@ -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
@ -21,7 +23,7 @@ from scrapy.core.downloader.contextfactory import (
_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,
@ -36,8 +38,8 @@ 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.interfaces import IListeningPort
from twisted.python.failure import Failure
from twisted.web.iweb import IBodyProducer
from scrapy.http import Response
@ -334,6 +336,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,

View File

@ -6,8 +6,8 @@ import re
import signal
import threading
from pathlib import Path
from typing import TYPE_CHECKING, Any, ClassVar
from unittest.mock import MagicMock
from typing import TYPE_CHECKING, Any, ClassVar, cast
from unittest.mock import MagicMock, patch
import pytest
@ -668,9 +668,10 @@ class TestAsyncCrawlerProcessReactorlessHelpers:
process, installed_handlers = self._bare_process(monkeypatch)
process._reactorless_loop = None
# No loop to schedule the shutdown task on, so it returns early, but it
# must still escalate the handler so a second signal forces a kill.
# must still escalate the handler so a second signal forces a fast
# shutdown.
process._signal_shutdown_reactorless(signal.SIGINT, None)
assert installed_handlers == [process._signal_kill_reactorless]
assert installed_handlers == [process._signal_fast_shutdown_reactorless]
def test_signal_kill_reactorless_without_loop(
self, monkeypatch: pytest.MonkeyPatch
@ -695,13 +696,13 @@ class TestAsyncCrawlerProcessReactorlessHelpers:
assert installed_handlers == [signal.SIG_IGN]
loop.call_soon_threadsafe.assert_not_called()
def test_shutdown_graceful_reactorless_main_task_already_done(
def test_shutdown_reactorless_main_task_already_done(
self, monkeypatch: pytest.MonkeyPatch
) -> None:
process, _ = self._bare_process(monkeypatch)
process._stop_after_crawl = False
async def noop() -> None:
async def noop(*args: Any, **kwargs: Any) -> None:
return None
monkeypatch.setattr(process, "stop", noop)
@ -714,25 +715,13 @@ class TestAsyncCrawlerProcessReactorlessHelpers:
main_task.set_result(None)
process._reactorless_main_task = main_task
# The main task is already done, so it is not cancelled.
loop.run_until_complete(process._shutdown_graceful_reactorless())
loop.run_until_complete(process._shutdown_reactorless(mode="graceful"))
assert not main_task.cancelled()
finally:
loop.close()
self._run_in_thread(run)
def test_create_shutdown_task_closed_loop(
self, monkeypatch: pytest.MonkeyPatch
) -> None:
process, _ = self._bare_process(monkeypatch)
loop = asyncio.new_event_loop()
loop.close()
process._reactorless_loop = loop
process._stop_after_crawl = True
# create_task() raises RuntimeError on a closed loop; the coroutine
# must be closed instead of leaking.
process._create_shutdown_task()
def test_cancel_all_tasks_logs_task_exception(self) -> None:
contexts: list[dict[str, Any]] = []
task_was_cancelled: list[bool] = []
@ -806,3 +795,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

View File

@ -239,10 +239,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 sleep(0.01)
p.kill(sig)
p.expect_exact("forcing unclean shutdown")
p.expect_exact("dropping downloader requests")
await 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()
@ -468,6 +471,9 @@ class TestAsyncCrawlerProcessSubprocess(TestCrawlerProcessSubprocessBase):
def test_reactorless_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_reactorless_shutdown_forced(self) -> None:
await self._test_shutdown_forced("reactorless_sleeping.py")

View File

@ -5,17 +5,24 @@ import logging
import subprocess
import sys
from typing import TYPE_CHECKING, Any
from unittest.mock import Mock
from unittest.mock import AsyncMock, Mock, patch
import pytest
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 Request
from scrapy.spiders import Spider
from scrapy.utils.defer import _schedule_coro, deferred_from_coro
from scrapy.utils.asyncio import sleep
from scrapy.utils.defer import (
_schedule_coro,
deferred_from_coro,
maybe_deferred_to_future,
)
from scrapy.utils.misc import build_from_crawler
from scrapy.utils.spider import DefaultSpider
from scrapy.utils.test import get_crawler
@ -136,6 +143,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 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):

View File

@ -1,6 +1,7 @@
from __future__ import annotations
from typing import TYPE_CHECKING, cast
from unittest.mock import patch
import pytest
from twisted.internet import defer
@ -151,3 +152,47 @@ async def test_exception_async_callback(
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(
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(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()