diff --git a/scrapy/core/engine.py b/scrapy/core/engine.py index d8a0f8a1b..7fb906eab 100644 --- a/scrapy/core/engine.py +++ b/scrapy/core/engine.py @@ -14,7 +14,7 @@ from time import time from traceback import format_exc from typing import TYPE_CHECKING, Any -from twisted.internet.defer import CancelledError, Deferred, inlineCallbacks, succeed +from twisted.internet.defer import CancelledError, Deferred, inlineCallbacks from twisted.python.failure import Failure from scrapy import signals @@ -27,10 +27,13 @@ from scrapy.exceptions import ( ScrapyDeprecationWarning, ) from scrapy.http import Request, Response -from scrapy.utils.asyncio import AsyncioLoopingCall, create_looping_call +from scrapy.utils.asyncio import ( + AsyncioLoopingCall, + create_looping_call, + is_asyncio_available, +) from scrapy.utils.defer import ( _schedule_coro, - deferred_f_from_coro_f, deferred_from_coro, maybe_deferred_to_future, ) @@ -79,10 +82,10 @@ class _Slot: self.inprogress.remove(request) self._maybe_fire_closing() - def close(self) -> Deferred[None]: + async def close(self) -> None: self.closing = Deferred() self._maybe_fire_closing() - return self.closing + await maybe_deferred_to_future(self.closing) def _maybe_fire_closing(self) -> None: if self.closing is not None and not self.inprogress: @@ -116,7 +119,9 @@ class ExecutionEngine: self.start_time: float | None = None self._start: AsyncIterator[Any] | None = None self._closewait: Deferred[None] | None = None - self._start_request_processing_dfd: Deferred[None] | None = None + self._start_request_processing_awaitable: ( + asyncio.Future[None] | Deferred[None] | None + ) = None downloader_cls: type[Downloader] = load_object(self.settings["DOWNLOADER"]) try: self.scheduler_cls: type[BaseScheduler] = self._get_scheduler_class( @@ -136,7 +141,8 @@ class ExecutionEngine: self.scraper: Scraper = Scraper(crawler) except Exception: - self.close() + if hasattr(self, "downloader"): + self.downloader.close() raise def _get_scheduler_class(self, settings: BaseSettings) -> type[BaseScheduler]: @@ -154,9 +160,15 @@ class ExecutionEngine: ScrapyDeprecationWarning, stacklevel=2, ) - return deferred_from_coro(self.start_async(_start_request_processing)) + return deferred_from_coro( + self.start_async(_start_request_processing=_start_request_processing) + ) - async def start_async(self, _start_request_processing: bool = True) -> None: + async def start_async(self, *, _start_request_processing: bool = True) -> None: + """Start the execution engine. + + .. versionadded:: VERSION + """ if self.running: raise RuntimeError("Engine already running") self.start_time = time() @@ -167,46 +179,63 @@ class ExecutionEngine: self.running = True self._closewait = Deferred() if _start_request_processing: - self._start_request_processing_dfd = self._start_request_processing() + coro = self._start_request_processing() + if is_asyncio_available(): + # not wrapping in a Deferred here to avoid https://github.com/twisted/twisted/issues/12470 + # (can happen when this is cancelled, e.g. in test_close_during_start_iteration()) + self._start_request_processing_awaitable = asyncio.ensure_future(coro) + else: + self._start_request_processing_awaitable = Deferred.fromCoroutine(coro) await maybe_deferred_to_future(self._closewait) def stop(self) -> Deferred[None]: - """Gracefully stop the execution engine""" + warnings.warn( + "ExecutionEngine.stop() is deprecated, use stop_async() instead", + ScrapyDeprecationWarning, + stacklevel=2, + ) + return deferred_from_coro(self.stop_async()) - @deferred_f_from_coro_f - async def _finish_stopping_engine(_: Any) -> None: - await self.signals.send_catch_log_async(signal=signals.engine_stopped) - if self._closewait: - self._closewait.callback(None) + async def stop_async(self) -> None: + """Gracefully stop the execution engine. + + .. versionadded:: VERSION + """ if not self.running: raise RuntimeError("Engine not running") self.running = False - if self._start_request_processing_dfd is not None: - self._start_request_processing_dfd.cancel() - self._start_request_processing_dfd = None - dfd = ( - self.close_spider(self.spider, reason="shutdown") - if self.spider is not None - else succeed(None) - ) - return dfd.addBoth(_finish_stopping_engine) + if self._start_request_processing_awaitable is not None: + 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.signals.send_catch_log_async(signal=signals.engine_stopped) + if self._closewait: + self._closewait.callback(None) def close(self) -> Deferred[None]: + warnings.warn( + "ExecutionEngine.close() is deprecated, use close_async() instead", + ScrapyDeprecationWarning, + stacklevel=2, + ) + return deferred_from_coro(self.close_async()) + + async def close_async(self) -> None: """ Gracefully close the execution engine. If it has already been started, stop it. In all cases, close the spider and the downloader. """ if self.running: - return self.stop() # will also close spider and downloader - if self.spider is not None: - return self.close_spider( - self.spider, reason="shutdown" + await self.stop_async() # will also close spider and downloader + elif self.spider is not None: + await self.close_spider_async( + reason="shutdown" ) # will also close downloader - if hasattr(self, "downloader"): + elif hasattr(self, "downloader"): self.downloader.close() - return succeed(None) def pause(self) -> None: self.paused = True @@ -242,7 +271,6 @@ class ExecutionEngine: ) self._slot.nextcall.schedule() - @deferred_f_from_coro_f async def _start_request_processing(self) -> None: """Starts consuming Spider.start() output and sending scheduled requests.""" @@ -266,13 +294,13 @@ class ExecutionEngine: return except Exception: # an error happened, log it and stop the engine - self._start_request_processing_dfd = None + self._start_request_processing_awaitable = None logger.error( "Error while processing requests from start()", exc_info=True, extra={"spider": self.spider}, ) - await maybe_deferred_to_future(self.stop()) + await self.stop_async() def _start_scheduled_requests(self) -> None: if self._slot is None or self._slot.closing is not None or self.paused: @@ -409,7 +437,12 @@ class ExecutionEngine: return deferred_from_coro(self.download_async(request)) async def download_async(self, request: Request) -> Response: - """Asynchronous version of download() that returns a Response.""" + """Return a coroutine which fires with a Response as result. + + Only downloader middlewares are applied. + + .. versionadded:: VERSION + """ if self.spider is None: raise RuntimeError(f"No open spider to crawl: {request}") try: @@ -481,7 +514,7 @@ class ExecutionEngine: self._start = await self.scraper.spidermw.process_start() if hasattr(scheduler, "open") and (d := scheduler.open(self.crawler.spider)): await maybe_deferred_to_future(d) - await maybe_deferred_to_future(self.scraper.open_spider()) + await self.scraper.open_spider_async() assert self.crawler.stats self.crawler.stats.open_spider(self.crawler.spider) await self.signals.send_catch_log_async( @@ -511,20 +544,33 @@ class ExecutionEngine: if self.spider_is_idle(): ex = detected_ex.get(CloseSpider, CloseSpider(reason="finished")) assert isinstance(ex, CloseSpider) # typing - self.close_spider(self.spider, reason=ex.reason) + _schedule_coro(self.close_spider_async(reason=ex.reason)) + + def close_spider(self, spider: Spider, reason: str = "cancelled") -> Deferred[None]: + 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)) + + async def close_spider_async(self, *, reason: str = "cancelled") -> None: + """Close (cancel) spider and clear all its outstanding requests. + + .. versionadded:: VERSION + """ + if self.spider is None: + raise RuntimeError("Spider not opened") - @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: - yield self._slot.closing + await maybe_deferred_to_future(self._slot.closing) return + spider = self.spider + logger.info( "Closing spider (%(reason)s)", {"reason": reason}, extra={"spider": spider} ) @@ -533,7 +579,7 @@ class ExecutionEngine: logger.error(msg, exc_info=True, extra={"spider": spider}) # noqa: LOG014 try: - yield self._slot.close() + await self._slot.close() except Exception: log_failure("Slot close failure") @@ -543,19 +589,19 @@ class ExecutionEngine: log_failure("Downloader close failure") try: - yield self.scraper.close_spider() + await self.scraper.close_spider_async() except Exception: log_failure("Scraper close failure") if hasattr(self._slot.scheduler, "close"): try: if (d := self._slot.scheduler.close(reason)) is not None: - yield d + await maybe_deferred_to_future(d) except Exception: log_failure("Scheduler close failure") try: - yield self.signals.send_catch_log_deferred( + await self.signals.send_catch_log_async( signal=signals.spider_closed, spider=spider, reason=reason, @@ -580,6 +626,6 @@ class ExecutionEngine: try: if (d := self._spider_closed_callback(spider)) is not None: - yield d + await maybe_deferred_to_future(d) except Exception: log_failure("Error running spider_closed_callback") diff --git a/scrapy/core/scraper.py b/scrapy/core/scraper.py index 90236bec5..a7fff5c94 100644 --- a/scrapy/core/scraper.py +++ b/scrapy/core/scraper.py @@ -24,6 +24,7 @@ from scrapy.http import Request, Response from scrapy.utils.asyncio import _parallel_asyncio, is_asyncio_available from scrapy.utils.defer import ( _defer_sleep_async, + _schedule_coro, aiter_errback, deferred_f_from_coro_f, deferred_from_coro, @@ -130,16 +131,19 @@ class Scraper: assert crawler.logformatter self.logformatter: LogFormatter = crawler.logformatter - @deferred_f_from_coro_f - async def open_spider(self, spider: Spider | None = None) -> None: - """Open the spider for scraping and allocate resources for it""" - if spider is not None: - warnings.warn( - "Passing a 'spider' argument to Scraper.open_spider() is deprecated.", - category=ScrapyDeprecationWarning, - stacklevel=2, - ) + def open_spider(self, spider: Spider | None = None) -> Deferred[None]: + warnings.warn( + "Scraper.open_spider() is deprecated, use open_spider_async() instead", + ScrapyDeprecationWarning, + stacklevel=2, + ) + return deferred_from_coro(self.open_spider_async()) + async def open_spider_async(self) -> None: + """Open the spider for scraping and allocate resources for it. + + .. versionadded:: VERSION + """ self.slot = Slot(self.crawler.settings.getint("SCRAPER_SLOT_MAX_ACTIVE_SIZE")) if not self.crawler.spider: raise RuntimeError( @@ -152,24 +156,30 @@ class Scraper: else: await maybe_deferred_to_future(self.itemproc.open_spider()) - def close_spider(self, spider: Spider | None = None) -> Deferred[list[None]]: - """Close the spider being scraped and release its resources""" - if spider is not None: - warnings.warn( - "Passing a 'spider' argument to Scraper.close_spider() is deprecated.", - category=ScrapyDeprecationWarning, - stacklevel=2, - ) + def close_spider(self, spider: Spider | None = None) -> Deferred[None]: + warnings.warn( + "Scraper.close_spider() is deprecated, use close_spider_async() instead", + ScrapyDeprecationWarning, + stacklevel=2, + ) + return deferred_from_coro(self.close_spider_async()) + async def close_spider_async(self) -> None: + """Close the spider being scraped and release its resources. + + .. versionadded:: VERSION + """ if self.slot is None: raise RuntimeError("Scraper slot not assigned") self.slot.closing = Deferred() - if self._itemproc_needs_spider["close_spider"]: - d = self.slot.closing.addCallback(self.itemproc.close_spider) - else: - d = self.slot.closing.addCallback(lambda _: self.itemproc.close_spider()) self._check_if_closing() - return d + await maybe_deferred_to_future(self.slot.closing) + if self._itemproc_needs_spider["close_spider"]: + await maybe_deferred_to_future( + self.itemproc.close_spider(self.crawler.spider) + ) + else: + await maybe_deferred_to_future(self.itemproc.close_spider()) def is_idle(self) -> bool: """Return True if there isn't any more spiders to process""" @@ -271,7 +281,10 @@ class Scraper: async def call_spider_async( self, result: Response | Failure, request: Request ) -> Iterable[Any] | AsyncIterator[Any]: - """Call the request callback or errback with the response or failure.""" + """Call the request callback or errback with the response or failure. + + .. versionadded:: 2.13 + """ await _defer_sleep_async() assert self.crawler.spider if isinstance(result, Response): @@ -315,8 +328,8 @@ class Scraper: exc = _failure.value if isinstance(exc, CloseSpider): assert self.crawler.engine is not None # typing - self.crawler.engine.close_spider( - self.crawler.spider, exc.reason or "cancelled" + _schedule_coro( + self.crawler.engine.close_spider_async(reason=exc.reason or "cancelled") ) return logkws = self.logformatter.spider_error( @@ -365,7 +378,10 @@ class Scraper: request: Request, response: Response | Failure, ) -> None: - """Pass items/requests produced by a callback to ``_process_spidermw_output()`` in parallel.""" + """Pass items/requests produced by a callback to ``_process_spidermw_output()`` in parallel. + + .. versionadded:: 2.13 + """ it: Iterable[_T] | AsyncIterator[_T] if is_asyncio_available(): if isinstance(result, AsyncIterator): @@ -444,6 +460,8 @@ class Scraper: *response* is the source of the item data. If the item does not come from response data, e.g. it was hard-coded, set it to ``None``. + + .. versionadded:: VERSION """ assert self.slot is not None # typing assert self.crawler.spider is not None # typing diff --git a/scrapy/crawler.py b/scrapy/crawler.py index 32cde0300..369dc2255 100644 --- a/scrapy/crawler.py +++ b/scrapy/crawler.py @@ -164,7 +164,7 @@ class Crawler: except Exception: self.crawling = False if self.engine is not None: - yield self.engine.close() + yield deferred_from_coro(self.engine.close_async()) raise async def crawl_async(self, *args: Any, **kwargs: Any) -> None: @@ -200,7 +200,7 @@ class Crawler: except Exception: self.crawling = False if self.engine is not None: - await deferred_to_future(self.engine.close()) + await self.engine.close_async() raise def _create_spider(self, *args: Any, **kwargs: Any) -> Spider: @@ -216,7 +216,8 @@ class Crawler: if self.crawling: self.crawling = False assert self.engine - yield self.engine.stop() + if self.engine.running: + yield deferred_from_coro(self.engine.stop_async()) async def stop_async(self) -> None: """Start a graceful stop of the crawler and complete when the crawler is stopped. diff --git a/scrapy/extensions/closespider.py b/scrapy/extensions/closespider.py index b4c6c73a0..a4362b182 100644 --- a/scrapy/extensions/closespider.py +++ b/scrapy/extensions/closespider.py @@ -18,6 +18,7 @@ from scrapy.utils.asyncio import ( call_later, create_looping_call, ) +from scrapy.utils.defer import _schedule_coro if TYPE_CHECKING: from twisted.internet.task import LoopingCall @@ -86,38 +87,31 @@ class CloseSpider: def error_count(self, failure: Failure, response: Response, spider: Spider) -> None: self.counter["errorcount"] += 1 if self.counter["errorcount"] == self.close_on["errorcount"]: - assert self.crawler.engine - self.crawler.engine.close_spider(spider, "closespider_errorcount") + self._close_spider("closespider_errorcount") def page_count(self, response: Response, request: Request, spider: Spider) -> None: self.counter["pagecount"] += 1 self.counter["pagecount_since_last_item"] += 1 if self.counter["pagecount"] == self.close_on["pagecount"]: - assert self.crawler.engine - self.crawler.engine.close_spider(spider, "closespider_pagecount") + self._close_spider("closespider_pagecount") return if self.close_on["pagecount_no_item"] and ( self.counter["pagecount_since_last_item"] >= self.close_on["pagecount_no_item"] ): - assert self.crawler.engine - self.crawler.engine.close_spider(spider, "closespider_pagecount_no_item") + self._close_spider("closespider_pagecount_no_item") def spider_opened(self, spider: Spider) -> None: assert self.crawler.engine self.task = call_later( - self.close_on["timeout"], - self.crawler.engine.close_spider, - spider, - "closespider_timeout", + self.close_on["timeout"], self._close_spider, "closespider_timeout" ) def item_scraped(self, item: Any, spider: Spider) -> None: self.counter["itemcount"] += 1 self.counter["pagecount_since_last_item"] = 0 if self.counter["itemcount"] == self.close_on["itemcount"]: - assert self.crawler.engine - self.crawler.engine.close_spider(spider, "closespider_itemcount") + self._close_spider("closespider_itemcount") def spider_closed(self, spider: Spider) -> None: if self.task: @@ -130,7 +124,7 @@ class CloseSpider: self.task_no_item = None def spider_opened_no_item(self, spider: Spider) -> None: - self.task_no_item = create_looping_call(self._count_items_produced, spider) + self.task_no_item = create_looping_call(self._count_items_produced) self.task_no_item.start(self.timeout_no_item, now=False) logger.info( @@ -141,7 +135,7 @@ class CloseSpider: def item_scraped_no_item(self, item: Any, spider: Spider) -> None: self.items_in_period += 1 - def _count_items_produced(self, spider: Spider) -> None: + def _count_items_produced(self) -> None: if self.items_in_period >= 1: self.items_in_period = 0 else: @@ -149,5 +143,8 @@ class CloseSpider: f"Closing spider since no items were produced in the last " f"{self.timeout_no_item} seconds." ) - assert self.crawler.engine - self.crawler.engine.close_spider(spider, "closespider_timeout_no_item") + self._close_spider("closespider_timeout_no_item") + + def _close_spider(self, reason: str) -> None: + assert self.crawler.engine + _schedule_coro(self.crawler.engine.close_spider_async(reason=reason)) diff --git a/scrapy/extensions/memusage.py b/scrapy/extensions/memusage.py index e425749f7..a345bbcc9 100644 --- a/scrapy/extensions/memusage.py +++ b/scrapy/extensions/memusage.py @@ -17,6 +17,7 @@ from scrapy import signals from scrapy.exceptions import NotConfigured from scrapy.mail import MailSender from scrapy.utils.asyncio import AsyncioLoopingCall, create_looping_call +from scrapy.utils.defer import _schedule_coro from scrapy.utils.engine import get_engine_status if TYPE_CHECKING: @@ -110,8 +111,8 @@ class MemoryUsage: self.crawler.stats.set_value("memusage/limit_notified", 1) if self.crawler.engine.spider is not None: - self.crawler.engine.close_spider( - self.crawler.engine.spider, "memusage_exceeded" + _schedule_coro( + self.crawler.engine.close_spider_async(reason="memusage_exceeded") ) else: self.crawler.stop() diff --git a/scrapy/shell.py b/scrapy/shell.py index 39366312f..4b2bdf6cf 100644 --- a/scrapy/shell.py +++ b/scrapy/shell.py @@ -25,7 +25,7 @@ from scrapy.spiders import Spider from scrapy.utils.conf import get_config from scrapy.utils.console import DEFAULT_PYTHON_SHELLS, start_python_console from scrapy.utils.datatypes import SequenceExclude -from scrapy.utils.defer import deferred_f_from_coro_f +from scrapy.utils.defer import _schedule_coro, deferred_f_from_coro_f from scrapy.utils.misc import load_object from scrapy.utils.reactor import is_asyncio_reactor_installed, set_asyncio_event_loop from scrapy.utils.response import open_in_browser @@ -127,7 +127,7 @@ class Shell: self.crawler.spider = spider assert self.crawler.engine await self.crawler.engine.open_spider_async(close_if_idle=False) - self.crawler.engine._start_request_processing() + _schedule_coro(self.crawler.engine._start_request_processing()) self.spider = spider def fetch( diff --git a/scrapy/signalmanager.py b/scrapy/signalmanager.py index 283060074..347eddfdb 100644 --- a/scrapy/signalmanager.py +++ b/scrapy/signalmanager.py @@ -77,6 +77,8 @@ class SignalManager: The keyword arguments are passed to the signal handlers (connected through the :meth:`connect` method). + + .. versionadded:: VERSION """ kwargs.setdefault("sender", self.sender) return await _signal.send_catch_log_async(signal, **kwargs) diff --git a/scrapy/utils/signal.py b/scrapy/utils/signal.py index 552fbaa90..1b890933b 100644 --- a/scrapy/utils/signal.py +++ b/scrapy/utils/signal.py @@ -120,6 +120,8 @@ async def send_catch_log_async( `. Returns a coroutine that completes once all signal handlers have finished. + + .. versionadded:: VERSION """ return await maybe_deferred_to_future( send_catch_log_deferred(signal, sender, *arguments, **named) diff --git a/tests/test_downloadermiddleware.py b/tests/test_downloadermiddleware.py index 6a0341414..2684cc1de 100644 --- a/tests/test_downloadermiddleware.py +++ b/tests/test_downloadermiddleware.py @@ -33,7 +33,7 @@ class TestManagerBase: crawler.engine = crawler._create_engine() await crawler.engine.open_spider_async() yield mwman - await maybe_deferred_to_future(crawler.engine.close_spider(crawler.spider)) + await crawler.engine.close_spider_async() @staticmethod async def _download( diff --git a/tests/test_engine.py b/tests/test_engine.py index 753f67f23..94afcc2ad 100644 --- a/tests/test_engine.py +++ b/tests/test_engine.py @@ -32,7 +32,6 @@ from scrapy.utils.defer import ( _schedule_coro, deferred_f_from_coro_f, deferred_from_coro, - deferred_to_future, maybe_deferred_to_future, ) from scrapy.utils.signal import disconnect_all @@ -417,10 +416,10 @@ class TestEngine(TestEngineBase): "reason": "custom_reason", } == run.signals_caught[signals.spider_closed] - @inlineCallbacks - def test_close_downloader(self): + @deferred_f_from_coro_f + async def test_close_downloader(self): e = ExecutionEngine(get_crawler(MySpider), lambda _: None) - yield e.close() + await e.close_async() def test_close_without_downloader(self): class CustomException(Exception): @@ -444,7 +443,7 @@ class TestEngine(TestEngineBase): _schedule_coro(e.start_async()) with pytest.raises(RuntimeError, match="Engine already running"): yield deferred_from_coro(e.start_async()) - yield e.stop() + yield deferred_from_coro(e.stop_async()) @pytest.mark.only_asyncio @deferred_f_from_coro_f @@ -455,7 +454,7 @@ class TestEngine(TestEngineBase): await e.open_spider_async() with pytest.raises(RuntimeError, match="Engine already running"): await asyncio.gather(e.start_async(), e.start_async()) - await deferred_to_future(e.stop()) + await e.stop_async() @inlineCallbacks def test_start_request_processing_exception(self): @@ -635,7 +634,7 @@ def test_request_scheduled_signal(caplog): class TestEngineCloseSpider: - """Tests for exception handling coverage during close_spider().""" + """Tests for exception handling coverage during close_spider_async().""" @pytest.fixture def crawler(self) -> Crawler: @@ -648,9 +647,14 @@ class TestEngineCloseSpider: 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)) + await engine.close_spider_async() + + @deferred_f_from_coro_f + async def test_no_spider(self, crawler: Crawler) -> None: + engine = ExecutionEngine(crawler, lambda _: None) + with pytest.raises(RuntimeError, match="Spider not opened"): + await engine.close_spider_async() @deferred_f_from_coro_f async def test_exception_slot( @@ -660,8 +664,7 @@ class TestEngineCloseSpider: 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)) + await engine.close_spider_async() assert "Slot close failure" in caplog.text @deferred_f_from_coro_f @@ -671,8 +674,7 @@ class TestEngineCloseSpider: 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)) + await engine.close_spider_async() assert "Downloader close failure" in caplog.text @deferred_f_from_coro_f @@ -682,8 +684,7 @@ class TestEngineCloseSpider: 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)) + await engine.close_spider_async() assert "Scraper close failure" in caplog.text @deferred_f_from_coro_f @@ -694,8 +695,7 @@ class TestEngineCloseSpider: 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)) + await engine.close_spider_async() assert "Scheduler close failure" in caplog.text @deferred_f_from_coro_f @@ -705,8 +705,7 @@ class TestEngineCloseSpider: 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)) + await engine.close_spider_async() assert "Error while sending spider_close signal" in caplog.text @deferred_f_from_coro_f @@ -716,8 +715,7 @@ class TestEngineCloseSpider: 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)) + await engine.close_spider_async() assert "Stats close failure" in caplog.text @deferred_f_from_coro_f @@ -726,6 +724,5 @@ class TestEngineCloseSpider: ) -> 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)) + await engine.close_spider_async() assert "Error running spider_closed_callback" in caplog.text diff --git a/tests/test_engine_loop.py b/tests/test_engine_loop.py index 6915e1463..ddf1f1fe0 100644 --- a/tests/test_engine_loop.py +++ b/tests/test_engine_loop.py @@ -4,7 +4,6 @@ from collections import deque from logging import ERROR from typing import TYPE_CHECKING -from testfixtures import LogCapture from twisted.internet.defer import Deferred from scrapy import Request, Spider, signals @@ -14,6 +13,8 @@ from tests.mockserver.http import MockServer from tests.test_scheduler import MemoryScheduler if TYPE_CHECKING: + import pytest + from scrapy.http import Response @@ -86,13 +87,15 @@ class TestMain: assert actual_urls == expected_urls, f"{actual_urls=} != {expected_urls=}" @deferred_f_from_coro_f - async def test_close_during_start_iteration(self): + async def test_close_during_start_iteration( + self, caplog: pytest.LogCaptureFixture + ) -> None: class TestSpider(Spider): name = "test" async def start(self): assert self.crawler.engine is not None - await maybe_deferred_to_future(self.crawler.engine.close()) + await self.crawler.engine.close_async() yield Request("data:,a") def parse(self, response): @@ -107,15 +110,14 @@ class TestMain: crawler = get_crawler(TestSpider, settings_dict=settings) crawler.signals.connect(track_url, signals.request_reached_downloader) - with LogCapture(level=ERROR) as log: + caplog.clear() + with caplog.at_level(ERROR): await maybe_deferred_to_future(crawler.crawl()) - assert len(log.records) == 1 - assert log.records[0].msg == "Error running spider_closed_callback" - finish_reason = crawler.stats.get_value("finish_reason") - assert finish_reason == "shutdown", f"{finish_reason=}" - expected_urls = [] - assert actual_urls == expected_urls, f"{actual_urls=} != {expected_urls=}" + assert not caplog.records + assert crawler.stats + assert crawler.stats.get_value("finish_reason") == "shutdown" + assert not actual_urls class TestRequestSendOrder: