This commit is contained in:
Adrian 2026-06-30 01:11:59 +08:00 committed by GitHub
commit da757c798a
No known key found for this signature in database
GPG Key ID: B5690EEEBB952194
13 changed files with 823 additions and 58 deletions

View File

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

View File

@ -13,6 +13,7 @@ concurrency:
jobs:
tests:
runs-on: macos-latest
timeout-minutes: 60
env:
PYTEST_ADDOPTS: -n auto
strategy:

View File

@ -13,6 +13,7 @@ concurrency:
jobs:
tests:
runs-on: ubuntu-latest
timeout-minutes: 60
env:
PYTEST_ADDOPTS: -n auto
strategy:

View File

@ -13,6 +13,7 @@ concurrency:
jobs:
tests:
runs-on: windows-latest
timeout-minutes: 60
env:
PYTEST_ADDOPTS: -n auto
strategy:

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,
)
@ -117,6 +117,8 @@ class Downloader:
DownloaderMiddlewareManager.from_crawler(crawler)
)
self._slot_gc_loop: AsyncioLoopingCall | LoopingCall | None = None
self._accepting_requests: bool = True
self._download_tasks: dict[Request, Deferred[None]] = {}
self.per_slot_settings: dict[str, dict[str, Any]] = self.settings.getdict(
"DOWNLOAD_SLOTS"
)
@ -176,6 +178,10 @@ class Downloader:
# passed as download_func into self.middleware.download() in self.fetch()
async def _enqueue_request(self, request: Request) -> Response:
if not self._accepting_requests:
raise DownloadCancelledError(
"The downloader is shutting down and not accepting new requests"
)
key, slot = self._get_slot(request)
request.meta[self.DOWNLOAD_SLOT] = key
slot.active.add(request)
@ -210,12 +216,21 @@ class Downloader:
while slot.queue and slot.free_transfer_slots() > 0:
slot.lastseen = now
request, queue_dfd = slot.queue.popleft()
_schedule_coro(self._wait_for_download(slot, request, queue_dfd))
download_dfd = deferred_from_coro(
self._wait_for_download(slot, request, queue_dfd)
)
assert isinstance(download_dfd, Deferred)
self._download_tasks[request] = download_dfd
download_dfd.addBoth(self._download_task_done, request)
# prevent burst if inter-request delays were configured
if delay:
self._process_queue(slot)
break
def _download_task_done(self, result: Any, request: Request) -> Any:
self._download_tasks.pop(request, None)
return result
def _latercall(self, slot: Slot) -> None:
slot.latercall = None
self._process_queue(slot)
@ -257,11 +272,48 @@ class Downloader:
try:
response = await self._download(slot, request)
except Exception:
queue_dfd.errback(Failure())
if not queue_dfd.called:
queue_dfd.errback(Failure())
else:
queue_dfd.callback(response) # awaited in _enqueue_request()
async def stop_async(self) -> int:
self._accepting_requests = False
dropped_count = 0
for slot in self.slots.values():
slot.close()
while slot.queue:
request, queue_dfd = slot.queue.popleft()
dropped_count += 1
if not queue_dfd.called:
queue_dfd.errback(
Failure(
DownloadCancelledError(
"Request dropped due to fast downloader shutdown"
)
)
)
self.signals.send_catch_log(
signal=signals.request_left_downloader,
request=request,
spider=self.crawler.spider,
)
for download_dfd in list(self._download_tasks.values()):
if download_dfd.called:
continue
dropped_count += 1
download_dfd.cancel()
if dropped_count:
await _defer_sleep_async()
return dropped_count
def close(self) -> None:
self._accepting_requests = False
self._stop_slot_gc()
for slot in self.slots.values():
slot.close()

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,
@ -119,6 +121,8 @@ class ExecutionEngine:
self.running: bool = False
self._starting: bool = False
self._stopping: bool = False
self._stop_mode: _StopMode = "graceful"
self._downloader_fast_stopped: bool = False
self.paused: bool = False
self._spider_closed_callback: Callable[
[Spider], Coroutine[Any, Any, None] | Deferred[None] | None
@ -202,23 +206,36 @@ class ExecutionEngine:
with contextlib.suppress(asyncio.exceptions.CancelledError):
await maybe_deferred_to_future(self._closewait)
def stop(self) -> Deferred[None]: # pragma: no cover
def stop(
self, *, mode: _StopMode = "graceful"
) -> Deferred[None]: # pragma: no cover
warnings.warn(
"ExecutionEngine.stop() is deprecated, use stop_async() instead",
ScrapyDeprecationWarning,
stacklevel=2,
)
return deferred_from_coro(self.stop_async())
return deferred_from_coro(self.stop_async(mode=mode))
async def stop_async(self) -> None:
async def stop_async(self, *, mode: _StopMode = "graceful") -> None:
"""Gracefully stop the execution engine.
.. versionadded:: 2.14
"""
if not self._starting:
mode = _normalize_stop_mode(mode, allow_force=False)
if not self._starting and not self._stopping:
raise RuntimeError("Engine not running")
self._stop_mode = _max_stop_mode(self._stop_mode, mode)
if self._stopping:
if self.spider is not None and self._stop_mode == "fast":
await self.close_spider_async(reason="shutdown", mode="fast")
if self._closewait:
await maybe_deferred_to_future(self._closewait)
return
self.running = self._starting = False
self._stopping = True
if self._start_request_processing_awaitable is not None:
@ -233,7 +250,7 @@ class ExecutionEngine:
self._start_request_processing_awaitable.cancel()
self._start_request_processing_awaitable = None
if self.spider is not None:
await self.close_spider_async(reason="shutdown")
await self.close_spider_async(reason="shutdown", mode=self._stop_mode)
await self.signals.send_catch_log_async(signal=signals.engine_stopped)
if self._closewait:
self._closewait.callback(None)
@ -404,6 +421,17 @@ class ExecutionEngine:
f"Incorrect type: expected Request, Response or Failure, got {type(result)}: {result!r}"
)
if (
isinstance(result, Failure)
and self._stop_mode == "fast"
and result.check(
DownloadCancelledError,
CancelledError,
asyncio.exceptions.CancelledError,
)
):
return
# downloader middleware can return requests (for example, redirects)
if isinstance(result, Request):
self.crawl(result)
@ -583,20 +611,51 @@ class ExecutionEngine:
_schedule_coro(self.close_spider_async(reason=ex.reason))
def close_spider(
self, spider: Spider, reason: str = "cancelled"
self,
spider: Spider,
reason: str = "cancelled",
mode: _StopMode = "graceful",
) -> Deferred[None]: # pragma: no cover
warnings.warn(
"ExecutionEngine.close_spider() is deprecated, use close_spider_async() instead",
ScrapyDeprecationWarning,
stacklevel=2,
)
return deferred_from_coro(self.close_spider_async(reason=reason))
return deferred_from_coro(self.close_spider_async(reason=reason, mode=mode))
async def close_spider_async(self, *, reason: str = "cancelled") -> None: # noqa: PLR0912
async def _fast_stop_downloader(self) -> None:
if self._downloader_fast_stopped:
return
self._downloader_fast_stopped = True
dropped_count = await self.downloader.stop_async()
assert self.crawler.stats
if dropped_count:
self.crawler.stats.inc_value(
"downloader/request_dropped_count", dropped_count
)
logger.info(
"Fast shutdown dropped %(count)d downloader requests",
{"count": dropped_count},
extra={"spider": self.spider},
)
# pylint: disable=too-many-statements
async def close_spider_async( # noqa: PLR0912, PLR0915
self,
*,
reason: str = "cancelled",
mode: _StopMode = "graceful",
) -> None:
"""Close (cancel) spider and clear all its outstanding requests.
.. versionadded:: 2.14
"""
mode = _normalize_stop_mode(mode, allow_force=False)
self._stop_mode = _max_stop_mode(self._stop_mode, mode)
if self.spider is None:
raise RuntimeError("Spider not opened")
@ -604,6 +663,8 @@ class ExecutionEngine:
raise RuntimeError("Engine slot not assigned")
if self._slot.closing is not None:
if self._stop_mode == "fast":
await self._fast_stop_downloader()
await maybe_deferred_to_future(self._slot.closing)
return
@ -613,6 +674,9 @@ class ExecutionEngine:
"Closing spider (%(reason)s)", {"reason": reason}, extra={"spider": spider}
)
if self._stop_mode == "fast":
await self._fast_stop_downloader()
try:
await self._slot.close()
except Exception:

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,
@ -42,7 +43,7 @@ from scrapy.utils.reactor import (
from scrapy.utils.reactorless import install_reactor_import_hook
if TYPE_CHECKING:
from collections.abc import Awaitable, Generator, Iterable
from collections.abc import Awaitable, Callable, Generator, Iterable
from scrapy.logformatter import LogFormatter
from scrapy.statscollectors import StatsCollector
@ -85,6 +86,15 @@ class Crawler:
self.request_fingerprinter: RequestFingerprinterProtocol | None = None
self.spider: Spider | None = None
self.engine: ExecutionEngine | None = None
self._force_stop_callback: (
Callable[[], Awaitable[None] | Deferred[None] | None] | None
) = None
def _set_force_stop_callback(
self,
callback: Callable[[], Awaitable[None] | Deferred[None] | None] | None,
) -> None:
self._force_stop_callback = callback
def _update_root_log_handler(self) -> None:
if get_scrapy_root_handler() is not None:
@ -233,7 +243,7 @@ class Crawler:
def _create_engine(self) -> ExecutionEngine:
return ExecutionEngine(self, lambda _: self.stop_async())
def stop(self) -> Deferred[None]:
def stop(self, *, mode: _StopMode = "graceful") -> Deferred[None]:
"""Start a graceful stop of the crawler and return a deferred that is
fired when the crawler is stopped."""
warnings.warn(
@ -241,18 +251,42 @@ class Crawler:
ScrapyDeprecationWarning,
stacklevel=2,
)
return deferred_from_coro(self.stop_async())
return deferred_from_coro(self.stop_async(mode=mode))
async def stop_async(self) -> None:
async def stop_async(self, *, mode: _StopMode = "graceful") -> None:
"""Start a graceful stop of the crawler and complete when the crawler is stopped.
.. versionadded:: 2.14
"""
if self.crawling:
self.crawling = False
assert self.engine
if self.engine.running:
await self.engine.stop_async()
mode = _normalize_stop_mode(mode)
was_crawling = self.crawling
self.crawling = False
if not was_crawling and mode == "graceful":
return
if mode == "force":
if self._force_stop_callback is not None:
await ensure_awaitable(self._force_stop_callback())
return
logger.warning(
"Force stop requested, but no process-level force stop callback is available. Falling back to fast stop."
)
mode = "fast"
if self.engine is None:
return
# During shutdown callbacks, graceful stop may be re-entered after
# the engine has already switched to non-running state.
if mode == "graceful" and not self.engine.running:
return
try:
await self.engine.stop_async(mode=mode)
except RuntimeError as exc:
if str(exc) != "Engine not running":
raise
@staticmethod
def _get_component(
@ -475,13 +509,16 @@ class CrawlerRunner(CrawlerRunnerBase):
self._active.discard(d)
self.bootstrap_failed |= not getattr(crawler, "spider", None) or failed
def stop(self) -> Deferred[Any]:
def stop(self, *, mode: _StopMode = "graceful") -> Deferred[Any]:
"""
Stops simultaneously all the crawling jobs taking place.
Returns a deferred that is fired when they all have ended.
"""
return DeferredList(deferred_from_coro(c.stop_async()) for c in self.crawlers)
mode = _normalize_stop_mode(mode)
return DeferredList(
deferred_from_coro(c.stop_async(mode=mode)) for c in self.crawlers
)
@inlineCallbacks
def join(self) -> Generator[Deferred[Any], Any, None]:
@ -599,15 +636,16 @@ class AsyncCrawlerRunner(CrawlerRunnerBase):
return task
async def stop(self) -> None:
async def stop(self, *, mode: _StopMode = "graceful") -> None:
"""
Stops simultaneously all the crawling jobs taking place.
Completes when they all have ended.
"""
mode = _normalize_stop_mode(mode)
if self.crawlers:
await asyncio.wait(
[asyncio.create_task(c.stop_async()) for c in self.crawlers]
[asyncio.create_task(c.stop_async(mode=mode)) for c in self.crawlers]
)
async def join(self) -> None:
@ -629,6 +667,14 @@ class CrawlerProcessBase(CrawlerRunnerBase):
configure_logging(self.settings, install_root_handler)
log_scrapy_info(self.settings)
def _create_crawler(self, spidercls: str | type[Spider]) -> Crawler:
crawler = super()._create_crawler(spidercls)
crawler._set_force_stop_callback(self._force_stop)
return crawler
def _force_stop(self) -> None:
self._stop_reactor()
@abstractmethod
def start(
self, stop_after_crawl: bool = True, install_signal_handlers: bool = True
@ -638,10 +684,17 @@ class CrawlerProcessBase(CrawlerRunnerBase):
def _signal_shutdown(self, signum: int, _: Any) -> None:
from twisted.internet import reactor
install_shutdown_handlers(self._signal_kill)
install_shutdown_handlers(self._signal_fast_shutdown)
self._log_shutdown(signum)
reactor.callFromThread(self._graceful_stop_reactor)
def _signal_fast_shutdown(self, signum: int, _: Any) -> None:
from twisted.internet import reactor
install_shutdown_handlers(self._signal_kill)
self._log_fast_shutdown(signum)
reactor.callFromThread(self._fast_stop_reactor)
def _signal_kill(self, signum: int, _: Any) -> None:
from twisted.internet import reactor
@ -653,7 +706,15 @@ class CrawlerProcessBase(CrawlerRunnerBase):
def _log_shutdown(signum: int) -> None:
signame = signal_names[signum]
logger.info(
"Received %(signame)s, shutting down gracefully. Send again to force ",
"Received %(signame)s, shutting down gracefully. Send again to stop faster.",
{"signame": signame},
)
@staticmethod
def _log_fast_shutdown(signum: int) -> None:
signame = signal_names[signum]
logger.info(
"Received %(signame)s twice, dropping downloader requests. Send again to force unclean shutdown",
{"signame": signame},
)
@ -661,7 +722,8 @@ class CrawlerProcessBase(CrawlerRunnerBase):
def _log_kill(signum: int) -> None:
signame = signal_names[signum]
logger.info(
"Received %(signame)s twice, forcing unclean shutdown", {"signame": signame}
"Received %(signame)s three times, forcing unclean shutdown",
{"signame": signame},
)
def _setup_reactor(self, install_signal_handlers: bool) -> None:
@ -703,13 +765,20 @@ class CrawlerProcessBase(CrawlerRunnerBase):
)
@abstractmethod
def _stop_dfd(self) -> Deferred[Any]:
def _stop_dfd(self, *, mode: _StopMode = "graceful") -> Deferred[Any]:
raise NotImplementedError
@inlineCallbacks
def _graceful_stop_reactor(self) -> Generator[Deferred[Any], Any, None]:
try:
yield self._stop_dfd()
yield self._stop_dfd(mode="graceful")
finally:
self._stop_reactor()
@inlineCallbacks
def _fast_stop_reactor(self) -> Generator[Deferred[Any], Any, None]:
try:
yield self._stop_dfd(mode="fast")
finally:
self._stop_reactor()
@ -762,10 +831,12 @@ class CrawlerProcess(CrawlerProcessBase, CrawlerRunner):
spidercls = self.spider_loader.load(spidercls)
init_reactor = not self._initialized_reactor
self._initialized_reactor = True
return Crawler(spidercls, self.settings, init_reactor=init_reactor)
crawler = Crawler(spidercls, self.settings, init_reactor=init_reactor)
crawler._set_force_stop_callback(self._force_stop)
return crawler
def _stop_dfd(self) -> Deferred[Any]:
return self.stop()
def _stop_dfd(self, *, mode: _StopMode = "graceful") -> Deferred[Any]:
return self.stop(mode=mode)
def start(
self, stop_after_crawl: bool = True, install_signal_handlers: bool = True
@ -862,8 +933,18 @@ class AsyncCrawlerProcess(CrawlerProcessBase, AsyncCrawlerRunner):
self._initialized_reactor = True
self._reactorless_main_task: asyncio.Future[None] | None = None
def _stop_dfd(self) -> Deferred[Any]:
return deferred_from_coro(self.stop())
def _force_stop(self) -> None:
if self.settings.getbool("TWISTED_REACTOR_ENABLED"):
self._stop_reactor()
return
if (loop := self._reactorless_loop) is None:
return
if (task := self._reactorless_main_task) is not None:
loop.call_soon_threadsafe(task.cancel)
def _stop_dfd(self, *, mode: _StopMode = "graceful") -> Deferred[Any]:
return deferred_from_coro(self.stop(mode=mode))
def start(
self, stop_after_crawl: bool = True, install_signal_handlers: bool = True
@ -1008,24 +1089,31 @@ class AsyncCrawlerProcess(CrawlerProcessBase, AsyncCrawlerRunner):
}
)
def _signal_shutdown_reactorless(self, signum: int, _: Any) -> None:
install_shutdown_handlers(self._signal_kill_reactorless)
self._log_shutdown(signum)
def _schedule_reactorless_shutdown(self, *, mode: _StopMode) -> None:
if (loop := self._reactorless_loop) is None:
return
loop.call_soon_threadsafe(self._create_shutdown_task)
def _create_shutdown_task() -> None:
coro = self._shutdown_reactorless(mode=mode)
try:
loop.create_task(coro)
except RuntimeError:
coro.close()
def _create_shutdown_task(self) -> None:
assert self._reactorless_loop
coro = self._shutdown_graceful_reactorless()
try:
self._reactorless_loop.create_task(coro)
except RuntimeError:
coro.close()
loop.call_soon_threadsafe(_create_shutdown_task)
async def _shutdown_graceful_reactorless(self) -> None:
await self.stop()
def _signal_shutdown_reactorless(self, signum: int, _: Any) -> None:
install_shutdown_handlers(self._signal_fast_shutdown_reactorless)
self._log_shutdown(signum)
self._schedule_reactorless_shutdown(mode="graceful")
def _signal_fast_shutdown_reactorless(self, signum: int, _: Any) -> None:
install_shutdown_handlers(self._signal_kill_reactorless)
self._log_fast_shutdown(signum)
self._schedule_reactorless_shutdown(mode="fast")
async def _shutdown_reactorless(self, *, mode: _StopMode) -> None:
await self.stop(mode=mode)
if not self._stop_after_crawl:
# wait until crawl tasks finish and cancel the future
await self.join()

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
@ -14,13 +16,14 @@ from twisted.web import server, static
from twisted.web.client import Agent, BrowserLikePolicyForHTTPS, readBody
from twisted.web.client import Response as TxResponse
from scrapy import Request
from scrapy.core.downloader import Downloader, Slot, tls
from scrapy.core.downloader.contextfactory import (
_load_context_factory_from_settings,
_ScrapyClientContextFactory,
)
from scrapy.core.downloader.handlers.http11 import _RequestBodyProducer
from scrapy.exceptions import ScrapyDeprecationWarning
from scrapy.exceptions import DownloadCancelledError, ScrapyDeprecationWarning
from scrapy.utils._deps_compat import (
PYOPENSSL_SET_CIPHER_LIST_TMP_CONN,
TWISTED_TLS_NEW_IMPL,
@ -35,7 +38,7 @@ from tests.mockserver.utils import ssl_context_factory
from tests.utils.decorators import coroutine_test
if TYPE_CHECKING:
from twisted.internet.defer import Deferred
from twisted.python.failure import Failure
from twisted.web.iweb import IBodyProducer
@ -310,6 +313,119 @@ async def test_fetch_deprecated_spider_arg():
await crawler.crawl_async()
@coroutine_test
async def test_stop_async_drops_queued_requests() -> None:
crawler = get_crawler(DefaultSpider)
crawler.spider = crawler._create_spider()
downloader = Downloader(crawler)
slot = Slot(concurrency=1, delay=0, randomize_delay=False)
downloader.slots["example.com"] = slot
request = Request("https://example.com")
queue_dfd: Deferred[Any] = Deferred()
failures: list[Failure] = []
queue_dfd.addErrback(failures.append)
slot.queue.append((request, queue_dfd))
dropped = await downloader.stop_async()
assert dropped == 1
assert len(failures) == 1
assert failures[0].check(DownloadCancelledError)
@coroutine_test
async def test_stop_async_rejects_new_requests() -> None:
crawler = get_crawler(DefaultSpider)
crawler.spider = crawler._create_spider()
downloader = Downloader(crawler)
await downloader.stop_async()
with pytest.raises(
DownloadCancelledError,
match="not accepting new requests",
):
await downloader._enqueue_request(Request("https://example.com"))
@coroutine_test
async def test_wait_for_download_errbacks_queue_deferred_on_error() -> None:
crawler = get_crawler(DefaultSpider)
downloader = Downloader(crawler)
slot = Slot(concurrency=1, delay=0, randomize_delay=False)
queue_dfd: Deferred[Any] = Deferred()
failures: list[Failure] = []
queue_dfd.addErrback(failures.append)
with patch.object(downloader, "_download", side_effect=RuntimeError("boom")):
await downloader._wait_for_download(
slot,
Request("https://example.com"),
queue_dfd,
)
assert len(failures) == 1
assert failures[0].check(RuntimeError)
@coroutine_test
async def test_wait_for_download_keeps_called_queue_deferred_on_error() -> None:
crawler = get_crawler(DefaultSpider)
downloader = Downloader(crawler)
slot = Slot(concurrency=1, delay=0, randomize_delay=False)
queue_dfd: Deferred[Any] = Deferred()
queue_dfd.callback(None)
with patch.object(downloader, "_download", side_effect=RuntimeError("boom")):
await downloader._wait_for_download(
slot,
Request("https://example.com"),
queue_dfd,
)
assert queue_dfd.called
@coroutine_test
async def test_stop_async_skips_called_queued_deferred() -> None:
crawler = get_crawler(DefaultSpider)
crawler.spider = crawler._create_spider()
downloader = Downloader(crawler)
slot = Slot(concurrency=1, delay=0, randomize_delay=False)
downloader.slots["example.com"] = slot
queue_dfd: Deferred[Any] = Deferred()
queue_dfd.callback(None)
slot.queue.append((Request("https://example.com"), queue_dfd))
dropped = await downloader.stop_async()
assert dropped == 1
@coroutine_test
async def test_stop_async_cancels_pending_download_tasks() -> None:
crawler = get_crawler(DefaultSpider)
downloader = Downloader(crawler)
done_dfd: Deferred[None] = Deferred()
done_dfd.callback(None)
pending_dfd: Deferred[None] = Deferred()
failures: list[Failure] = []
pending_dfd.addErrback(failures.append)
downloader._download_tasks[Request("https://done.example")] = done_dfd
downloader._download_tasks[Request("https://pending.example")] = pending_dfd
dropped = await downloader.stop_async()
assert dropped == 1
assert len(failures) == 1
assert failures[0].check(CancelledError)
def test_deprecated_tls_module_names() -> None:
with pytest.warns(
ScrapyDeprecationWarning,

View File

@ -4,7 +4,8 @@ import asyncio
import logging
import re
from pathlib import Path
from typing import Any, ClassVar
from typing import Any, ClassVar, cast
from unittest.mock import patch
import pytest
from zope.interface.exceptions import MultipleInvalid
@ -826,3 +827,282 @@ async def test_deprecated_crawler_stop() -> None:
ScrapyDeprecationWarning, match=r"Crawler.stop\(\) is deprecated"
):
await maybe_deferred_to_future(crawler.stop())
@coroutine_test
async def test_crawler_stop_async_invalid_mode() -> None:
crawler = get_crawler(DefaultSpider)
with pytest.raises(ValueError, match=r"Unknown stop mode"):
await crawler.stop_async(mode="invalid") # type: ignore[arg-type]
@coroutine_test
async def test_crawler_graceful_stop_non_running_engine_is_noop() -> None:
crawler = get_crawler(DefaultSpider)
crawler.crawling = True
class DummyEngine:
running = False
called = False
async def stop_async(self, *, mode: str = "graceful") -> None:
self.called = True
dummy_engine = DummyEngine()
crawler.engine = dummy_engine # type: ignore[assignment]
await crawler.stop_async(mode="graceful")
assert dummy_engine.called is False
@coroutine_test
async def test_crawler_force_stop_falls_back_to_fast(
caplog: pytest.LogCaptureFixture,
) -> None:
crawler = get_crawler(DefaultSpider)
class DummyEngine:
called_mode: str | None = None
async def stop_async(self, *, mode: str = "graceful") -> None:
self.called_mode = mode
dummy_engine = DummyEngine()
crawler.engine = dummy_engine # type: ignore[assignment]
with caplog.at_level(logging.WARNING):
await crawler.stop_async(mode="force")
assert dummy_engine.called_mode == "fast"
assert "Falling back to fast stop" in caplog.text
@coroutine_test
async def test_crawler_force_stop_uses_force_callback() -> None:
crawler = get_crawler(DefaultSpider)
called = False
def force_stop_callback() -> None:
nonlocal called
called = True
crawler._set_force_stop_callback(force_stop_callback)
await crawler.stop_async(mode="force")
assert called
@coroutine_test
async def test_crawler_stop_async_without_engine_is_noop() -> None:
crawler = get_crawler(DefaultSpider)
crawler.crawling = True
await crawler.stop_async(mode="graceful")
assert crawler.crawling is False
@coroutine_test
async def test_crawler_stop_async_ignores_engine_not_running_runtime_error() -> None:
crawler = get_crawler(DefaultSpider)
crawler.crawling = True
class DummyEngine:
running = True
called = False
async def stop_async(self, *, mode: str = "graceful") -> None:
self.called = True
raise RuntimeError("Engine not running")
dummy_engine = DummyEngine()
crawler.engine = dummy_engine # type: ignore[assignment]
await crawler.stop_async(mode="graceful")
assert dummy_engine.called is True
@coroutine_test
async def test_crawler_stop_async_reraises_other_runtime_errors() -> None:
crawler = get_crawler(DefaultSpider)
crawler.crawling = True
class DummyEngine:
running = True
async def stop_async(self, *, mode: str = "graceful") -> None:
raise RuntimeError("different runtime error")
crawler.engine = DummyEngine() # type: ignore[assignment]
with pytest.raises(RuntimeError, match="different runtime error"):
await crawler.stop_async(mode="graceful")
@pytest.mark.requires_reactor
@coroutine_test
async def test_crawler_process_force_stop_via_public_crawler_api() -> None:
crawler_process = CrawlerProcess(install_root_handler=False)
called = False
def stop_reactor(_: Any = None) -> None:
nonlocal called
called = True
crawler_process._stop_reactor = stop_reactor # type: ignore[method-assign]
crawler = crawler_process.create_crawler(DefaultSpider)
await crawler.stop_async(mode="force")
assert called
@pytest.mark.only_asyncio
@pytest.mark.requires_reactor
@coroutine_test
async def test_async_crawler_process_force_stop_reactor_enabled_via_public_crawler_api() -> (
None
):
crawler_process = AsyncCrawlerProcess(
{"TWISTED_REACTOR_ENABLED": True},
install_root_handler=False,
)
called = False
def stop_reactor(_: Any = None) -> None:
nonlocal called
called = True
crawler_process._stop_reactor = stop_reactor # type: ignore[method-assign]
crawler = crawler_process.create_crawler(DefaultSpider)
await crawler.stop_async(mode="force")
assert called
@pytest.mark.only_asyncio
@coroutine_test
async def test_async_crawler_process_force_stop_reactorless_without_main_task(
reactor_pytest: str,
) -> None:
if reactor_pytest != "none":
pytest.skip("This test is only for --reactor=none")
crawler_process = AsyncCrawlerProcess(
{"TWISTED_REACTOR_ENABLED": False},
install_root_handler=False,
)
assert crawler_process._reactorless_loop is not None
assert crawler_process._reactorless_main_task is None
crawler = crawler_process.create_crawler(DefaultSpider)
await crawler.stop_async(mode="force")
@pytest.mark.only_asyncio
@coroutine_test
async def test_async_crawler_process_force_stop_reactorless_without_loop(
reactor_pytest: str,
) -> None:
if reactor_pytest != "none":
pytest.skip("This test is only for --reactor=none")
crawler_process = AsyncCrawlerProcess(
{"TWISTED_REACTOR_ENABLED": False},
install_root_handler=False,
)
crawler_process._reactorless_loop = None
crawler_process._reactorless_main_task = None
crawler = crawler_process.create_crawler(DefaultSpider)
await crawler.stop_async(mode="force")
@pytest.mark.only_asyncio
@coroutine_test
async def test_async_crawler_process_force_stop_reactorless_with_task(
reactor_pytest: str,
) -> None:
if reactor_pytest != "none":
pytest.skip("This test is only for --reactor=none")
crawler_process = AsyncCrawlerProcess(
{"TWISTED_REACTOR_ENABLED": False},
install_root_handler=False,
)
class DummyLoop:
callback = None
def call_soon_threadsafe(self, callback) -> None:
self.callback = callback
class DummyTask:
called = False
def cancel(self) -> None:
self.called = True
loop = DummyLoop()
task = DummyTask()
crawler_process._reactorless_loop = cast("asyncio.AbstractEventLoop", loop)
crawler_process._reactorless_main_task = cast("asyncio.Future[None]", task)
crawler = crawler_process.create_crawler(DefaultSpider)
await crawler.stop_async(mode="force")
assert loop.callback is not None
loop.callback()
assert task.called is True
def test_async_crawler_process_schedule_reactorless_shutdown_without_loop() -> None:
crawler_process = object.__new__(AsyncCrawlerProcess)
crawler_process._reactorless_loop = None
crawler_process._schedule_reactorless_shutdown(mode="graceful")
def test_async_crawler_process_schedule_reactorless_shutdown_runtime_error() -> None:
crawler_process = object.__new__(AsyncCrawlerProcess)
class DummyLoop:
scheduled = False
create_task_called = False
def call_soon_threadsafe(self, callback) -> None:
self.scheduled = True
callback()
def create_task(self, coro) -> None:
self.create_task_called = True
raise RuntimeError("event loop is closing")
class DummyCoro:
closed = False
def close(self) -> None:
self.closed = True
loop = DummyLoop()
coro = DummyCoro()
called_mode: str | None = None
def shutdown_reactorless(*, mode: str) -> DummyCoro:
nonlocal called_mode
called_mode = mode
return coro
crawler_process._reactorless_loop = cast("asyncio.AbstractEventLoop", loop)
with patch.object(
crawler_process, "_shutdown_reactorless", new=shutdown_reactorless
):
crawler_process._schedule_reactorless_shutdown(mode="graceful")
assert called_mode == "graceful"
assert loop.scheduled
assert loop.create_task_called
assert coro.closed

View File

@ -231,10 +231,13 @@ class TestCrawlerProcessSubprocessBase(ScriptRunnerMixin):
p.expect_exact("Crawled (200)")
p.kill(sig)
p.expect_exact("shutting down gracefully")
# sending the second signal too fast often causes problems
# sending a new signal too fast often causes problems
await async_sleep(0.01)
p.kill(sig)
p.expect_exact("forcing unclean shutdown")
p.expect_exact("dropping downloader requests")
await async_sleep(0.01)
p.kill(sig)
p.expect_exact("forcing unclean shutdown", timeout=20)
p.wait() # type: ignore[no-untyped-call]
if p.proc.stdin:
p.proc.stdin.close()
@ -425,6 +428,9 @@ class TestAsyncCrawlerProcessSubprocess(TestCrawlerProcessSubprocessBase):
def test_shutdown_graceful(self) -> None:
self._test_shutdown_graceful("reactorless_sleeping.py")
def test_shutdown_graceful_stop_after_crawl_false(self) -> None:
self._test_shutdown_graceful("reactorless_sleeping_no_stop_after_crawl.py")
@coroutine_test
async def test_shutdown_forced(self) -> None:
await self._test_shutdown_forced("reactorless_sleeping.py")

View File

@ -7,7 +7,7 @@ import sys
from collections import defaultdict
from dataclasses import dataclass
from typing import TYPE_CHECKING, Any, cast
from unittest.mock import Mock, call
from unittest.mock import AsyncMock, Mock, call, patch
from urllib.parse import urlparse
import attr
@ -16,11 +16,12 @@ from itemadapter import ItemAdapter
from pydispatch import dispatcher
from testfixtures import LogCapture
from twisted.internet import defer
from twisted.python.failure import Failure
from scrapy import signals
from scrapy.core.engine import ExecutionEngine, _Slot
from scrapy.core.scheduler import BaseScheduler
from scrapy.exceptions import CloseSpider, IgnoreRequest
from scrapy.exceptions import CloseSpider, DownloadCancelledError, IgnoreRequest
from scrapy.http import Headers, Request, Response
from scrapy.item import Field, Item
from scrapy.linkextractors import LinkExtractor
@ -35,11 +36,10 @@ from scrapy.utils.signal import disconnect_all
from scrapy.utils.spider import DefaultSpider
from scrapy.utils.test import get_crawler
from tests import get_testdata
from tests.utils import async_sleep
from tests.utils.decorators import coroutine_test, inline_callbacks_test
if TYPE_CHECKING:
from twisted.python.failure import Failure
from scrapy.core.scheduler import Scheduler
from scrapy.crawler import Crawler
from tests.mockserver.http import MockServer
@ -451,6 +451,71 @@ class TestEngine(TestEngineBase):
yield deferred_from_coro(e.start_async())
yield deferred_from_coro(e.stop_async())
@coroutine_test
async def test_stop_async_force_mode_not_supported(self) -> None:
engine = ExecutionEngine(get_crawler(DefaultSpider), lambda _: None)
with pytest.raises(ValueError, match="force stop mode is not supported"):
await engine.stop_async(mode="force")
@coroutine_test
async def test_stop_async_not_running_raises(self) -> None:
engine = ExecutionEngine(get_crawler(DefaultSpider), lambda _: None)
with pytest.raises(RuntimeError, match="Engine not running"):
await engine.stop_async()
@coroutine_test
async def test_stop_async_reentrant_fast_waits_for_closewait(self) -> None:
engine = ExecutionEngine(get_crawler(DefaultSpider), lambda _: None)
engine.spider = Mock()
engine._stopping = True
engine._closewait = defer.Deferred()
with patch.object(
engine, "close_spider_async", new_callable=AsyncMock
) as close:
stop_dfd = deferred_from_coro(engine.stop_async(mode="fast"))
await async_sleep(0)
close.assert_called_once_with(reason="shutdown", mode="fast")
assert not stop_dfd.called
assert engine._closewait
engine._closewait.callback(None)
await maybe_deferred_to_future(stop_dfd)
@coroutine_test
async def test_stop_async_reentrant_graceful_without_spider_or_closewait(
self,
) -> None:
engine = ExecutionEngine(get_crawler(DefaultSpider), lambda _: None)
engine._stopping = True
with patch.object(
engine, "close_spider_async", new_callable=AsyncMock
) as close:
await engine.stop_async(mode="graceful")
close.assert_not_called()
@coroutine_test
async def test_handle_downloader_output_ignores_fast_cancelled_failures(
self,
) -> None:
engine = ExecutionEngine(get_crawler(DefaultSpider), lambda _: None)
engine.spider = Mock()
engine._stop_mode = "fast"
enqueue_scrape = Mock()
engine.scraper.enqueue_scrape = enqueue_scrape # type: ignore[method-assign]
result = Failure(DownloadCancelledError("dropped during fast stop"))
await maybe_deferred_to_future(
engine._handle_downloader_output(result, Request("https://example.com"))
)
enqueue_scrape.assert_not_called()
@pytest.mark.only_asyncio
@coroutine_test
async def test_start_already_running_exception_asyncio(self):
@ -767,3 +832,45 @@ class TestEngineCloseSpider:
await engine.open_spider_async()
await engine.close_spider_async()
assert "Error running spider_closed_callback" in caplog.text
@coroutine_test
async def test_fast_close_stops_downloader_and_records_dropped_requests(
self, crawler: Crawler
) -> None:
engine = ExecutionEngine(crawler, lambda _: None)
crawler.engine = engine
await engine.open_spider_async()
calls = 0
async def fast_stop_downloader() -> int:
nonlocal calls
calls += 1
return 3
with patch.object(engine.downloader, "stop_async", fast_stop_downloader):
await engine.close_spider_async(mode="fast")
assert calls == 1
assert crawler.stats
assert crawler.stats.get_value("downloader/request_dropped_count") == 3
@coroutine_test
async def test_fast_stop_downloader_is_idempotent(self, crawler: Crawler) -> None:
engine = ExecutionEngine(crawler, lambda _: None)
crawler.engine = engine
await engine.open_spider_async()
calls = 0
async def fast_stop_downloader() -> int:
nonlocal calls
calls += 1
return 1
with patch.object(engine.downloader, "stop_async", fast_stop_downloader):
await engine._fast_stop_downloader()
await engine._fast_stop_downloader()
assert calls == 1
await engine.close_spider_async()