Rewrite ExecutionEngine.close_spider() to inlineCallbacks. (#6986)

This commit is contained in:
Andrey Rakhmatullin 2025-08-06 11:01:55 +05:00 committed by GitHub
parent baa579df62
commit a3daa3612e
No known key found for this signature in database
GPG Key ID: B5690EEEBB952194
2 changed files with 150 additions and 59 deletions

View File

@ -12,7 +12,7 @@ import logging
import warnings
from time import time
from traceback import format_exc
from typing import TYPE_CHECKING, Any, cast
from typing import TYPE_CHECKING, Any
from twisted.internet.defer import CancelledError, Deferred, inlineCallbacks, succeed
from twisted.python.failure import Failure
@ -513,75 +513,73 @@ class ExecutionEngine:
assert isinstance(ex, CloseSpider) # typing
self.close_spider(self.spider, reason=ex.reason)
def close_spider(self, spider: Spider, reason: str = "cancelled") -> Deferred[None]:
@inlineCallbacks
def close_spider(
self, spider: Spider, reason: str = "cancelled"
) -> Generator[Deferred[Any], Any, None]:
"""Close (cancel) spider and clear all its outstanding requests"""
if self._slot is None:
raise RuntimeError("Engine slot not assigned")
if self._slot.closing is not None:
return self._slot.closing
yield self._slot.closing
return
logger.info(
"Closing spider (%(reason)s)", {"reason": reason}, extra={"spider": spider}
)
dfd = self._slot.close()
def log_failure(msg: str) -> None:
logger.error(msg, exc_info=True, extra={"spider": spider}) # noqa: LOG014
def log_failure(msg: str) -> Callable[[Failure], None]:
def errback(failure: Failure) -> None:
logger.error(
msg, exc_info=failure_to_exc_info(failure), extra={"spider": spider}
)
try:
yield self._slot.close()
except Exception:
log_failure("Slot close failure")
return errback
try:
self.downloader.close()
except Exception:
log_failure("Downloader close failure")
dfd.addBoth(lambda _: self.downloader.close())
dfd.addErrback(log_failure("Downloader close failure"))
dfd.addBoth(lambda _: self.scraper.close_spider())
dfd.addErrback(log_failure("Scraper close failure"))
try:
yield self.scraper.close_spider()
except Exception:
log_failure("Scraper close failure")
if hasattr(self._slot.scheduler, "close"):
dfd.addBoth(lambda _: cast("_Slot", self._slot).scheduler.close(reason))
dfd.addErrback(log_failure("Scheduler close failure"))
try:
if (d := self._slot.scheduler.close(reason)) is not None:
yield d
except Exception:
log_failure("Scheduler close failure")
dfd.addBoth(
lambda _: self.signals.send_catch_log_deferred(
try:
yield self.signals.send_catch_log_deferred(
signal=signals.spider_closed,
spider=spider,
reason=reason,
)
)
dfd.addErrback(log_failure("Error while sending spider_close signal"))
except Exception:
log_failure("Error while sending spider_close signal")
def close_stats(_: Any) -> None:
assert self.crawler.stats
assert self.crawler.stats
try:
self.crawler.stats.close_spider(spider, reason=reason)
except Exception:
log_failure("Stats close failure")
dfd.addBoth(close_stats)
dfd.addErrback(log_failure("Stats close failure"))
dfd.addBoth(
lambda _: logger.info(
"Spider closed (%(reason)s)",
{"reason": reason},
extra={"spider": spider},
)
logger.info(
"Spider closed (%(reason)s)",
{"reason": reason},
extra={"spider": spider},
)
def unassign_slot(_: Any) -> None:
self._slot = None
self._slot = None
self.spider = None
dfd.addBoth(unassign_slot)
dfd.addErrback(log_failure("Error while unassigning slot"))
def unassign_spider(_: Any) -> None:
self.spider = None
dfd.addBoth(unassign_spider)
dfd.addErrback(log_failure("Error while unassigning spider"))
dfd.addBoth(lambda _: self._spider_closed_callback(spider))
dfd.addErrback(log_failure("Error running spider_closed_callback"))
return dfd
try:
if (d := self._spider_closed_callback(spider)) is not None:
yield d
except Exception:
log_failure("Error running spider_closed_callback")

View File

@ -1,14 +1,4 @@
"""
Scrapy engine tests
This starts a testing web server (using twisted.server.Site) and then crawls it
with the Scrapy crawler.
To view the testing web server in a browser you can start it by running this
module with the ``runserver`` argument::
python test_engine.py runserver
"""
from __future__ import annotations
import asyncio
import re
@ -17,6 +7,7 @@ import sys
from collections import defaultdict
from dataclasses import dataclass
from logging import DEBUG
from typing import TYPE_CHECKING, cast
from unittest.mock import Mock, call
from urllib.parse import urlparse
@ -48,7 +39,12 @@ 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.mockserver.http import MockServer
if TYPE_CHECKING:
from scrapy.core.scheduler import Scheduler
from scrapy.crawler import Crawler
from scrapy.statscollectors import MemoryStatsCollector
from tests.mockserver.http import MockServer
class MyItem(Item):
@ -636,3 +632,100 @@ def test_request_scheduled_signal(caplog):
f"{scheduler.enqueued!r} != [{keep_request!r}]"
)
crawler.signals.disconnect(signal_handler, request_scheduled)
class TestEngineCloseSpider:
"""Tests for exception handling coverage during close_spider()."""
@pytest.fixture
def crawler(self) -> Crawler:
crawler = get_crawler(DefaultSpider)
crawler.spider = crawler._create_spider()
return crawler
@deferred_f_from_coro_f
async def test_no_slot(self, crawler: Crawler) -> None:
engine = ExecutionEngine(crawler, lambda _: None)
await engine.open_spider_async()
engine._slot = None
assert crawler.spider
with pytest.raises(RuntimeError, match="Engine slot not assigned"):
await maybe_deferred_to_future(engine.close_spider(crawler.spider))
@deferred_f_from_coro_f
async def test_exception_slot(
self, crawler: Crawler, caplog: pytest.LogCaptureFixture
) -> None:
engine = ExecutionEngine(crawler, lambda _: None)
await engine.open_spider_async()
assert engine._slot
del engine._slot.heartbeat
assert crawler.spider
await maybe_deferred_to_future(engine.close_spider(crawler.spider))
assert "Slot close failure" in caplog.text
@deferred_f_from_coro_f
async def test_exception_downloader(
self, crawler: Crawler, caplog: pytest.LogCaptureFixture
) -> None:
engine = ExecutionEngine(crawler, lambda _: None)
await engine.open_spider_async()
del engine.downloader.slots
assert crawler.spider
await maybe_deferred_to_future(engine.close_spider(crawler.spider))
assert "Downloader close failure" in caplog.text
@deferred_f_from_coro_f
async def test_exception_scraper(
self, crawler: Crawler, caplog: pytest.LogCaptureFixture
) -> None:
engine = ExecutionEngine(crawler, lambda _: None)
await engine.open_spider_async()
engine.scraper.slot = None
assert crawler.spider
await maybe_deferred_to_future(engine.close_spider(crawler.spider))
assert "Scraper close failure" in caplog.text
@deferred_f_from_coro_f
async def test_exception_scheduler(
self, crawler: Crawler, caplog: pytest.LogCaptureFixture
) -> None:
engine = ExecutionEngine(crawler, lambda _: None)
await engine.open_spider_async()
assert engine._slot
del cast("Scheduler", engine._slot.scheduler).dqs
assert crawler.spider
await maybe_deferred_to_future(engine.close_spider(crawler.spider))
assert "Scheduler close failure" in caplog.text
@deferred_f_from_coro_f
async def test_exception_signal(
self, crawler: Crawler, caplog: pytest.LogCaptureFixture
) -> None:
engine = ExecutionEngine(crawler, lambda _: None)
await engine.open_spider_async()
del engine.signals
assert crawler.spider
await maybe_deferred_to_future(engine.close_spider(crawler.spider))
assert "Error while sending spider_close signal" in caplog.text
@deferred_f_from_coro_f
async def test_exception_stats(
self, crawler: Crawler, caplog: pytest.LogCaptureFixture
) -> None:
engine = ExecutionEngine(crawler, lambda _: None)
await engine.open_spider_async()
del cast("MemoryStatsCollector", crawler.stats).spider_stats
assert crawler.spider
await maybe_deferred_to_future(engine.close_spider(crawler.spider))
assert "Stats close failure" in caplog.text
@deferred_f_from_coro_f
async def test_exception_callback(
self, crawler: Crawler, caplog: pytest.LogCaptureFixture
) -> None:
engine = ExecutionEngine(crawler, lambda _: defer.fail(ValueError()))
await engine.open_spider_async()
assert crawler.spider
await maybe_deferred_to_future(engine.close_spider(crawler.spider))
assert "Error running spider_closed_callback" in caplog.text