mirror of https://github.com/scrapy/scrapy.git
Support fast crawler stops
This commit is contained in:
parent
320e40a044
commit
8613d390d7
|
|
@ -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()
|
||||
|
|
|
|||
|
|
@ -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:
|
||||
|
|
|
|||
|
|
@ -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()
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
@ -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"))
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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")
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
Loading…
Reference in New Issue