diff --git a/scrapy/commands/shell.py b/scrapy/commands/shell.py index 9ca383965..080f62382 100644 --- a/scrapy/commands/shell.py +++ b/scrapy/commands/shell.py @@ -12,6 +12,7 @@ from typing import TYPE_CHECKING, Any from scrapy.commands import ScrapyCommand from scrapy.http import Request from scrapy.shell import Shell +from scrapy.utils.defer import _schedule_coro from scrapy.utils.spider import DefaultSpider, spidercls_for_request from scrapy.utils.url import guess_scheme @@ -84,7 +85,7 @@ class Command(ScrapyCommand): crawler._apply_settings() # The Shell class needs a persistent engine in the crawler crawler.engine = crawler._create_engine() - crawler.engine.start(_start_request_processing=False) + _schedule_coro(crawler.engine.start_async(_start_request_processing=False)) self._start_crawler_thread() diff --git a/scrapy/core/engine.py b/scrapy/core/engine.py index 669005e32..e6ce58376 100644 --- a/scrapy/core/engine.py +++ b/scrapy/core/engine.py @@ -29,6 +29,7 @@ from scrapy.exceptions import ( from scrapy.http import Request, Response from scrapy.utils.asyncio import AsyncioLoopingCall, create_looping_call from scrapy.utils.defer import ( + _schedule_coro, deferred_f_from_coro_f, deferred_from_coro, maybe_deferred_to_future, @@ -135,9 +136,14 @@ class ExecutionEngine: return scheduler_cls def start(self, _start_request_processing=True) -> Deferred[None]: + warnings.warn( + "ExecutionEngine.start() is deprecated, use start_async() instead", + ScrapyDeprecationWarning, + stacklevel=2, + ) return deferred_from_coro(self.start_async(_start_request_processing)) - async def start_async(self, _start_request_processing=True) -> None: + async def start_async(self, _start_request_processing: bool = True) -> None: if self.running: raise RuntimeError("Engine already running") self.start_time = time() @@ -218,7 +224,9 @@ class ExecutionEngine: if isinstance(item_or_request, Request): self.crawl(item_or_request) else: - self.scraper.start_itemproc(item_or_request, response=None) + _schedule_coro( + self.scraper.start_itemproc_async(item_or_request, response=None) + ) self._slot.nextcall.schedule() @deferred_f_from_coro_f @@ -438,6 +446,11 @@ class ExecutionEngine: self._slot.nextcall.schedule() def open_spider(self, spider: Spider, close_if_idle: bool = True) -> Deferred[None]: + warnings.warn( + "ExecutionEngine.open_spider() is deprecated, use open_spider_async() instead", + ScrapyDeprecationWarning, + stacklevel=2, + ) return deferred_from_coro( self.open_spider_async(spider, close_if_idle=close_if_idle) ) diff --git a/scrapy/core/scraper.py b/scrapy/core/scraper.py index 6b80ba9bf..47d9ae2e2 100644 --- a/scrapy/core/scraper.py +++ b/scrapy/core/scraper.py @@ -193,7 +193,7 @@ class Scraper: try: # call the spider middlewares and the request callback with the response output = await self.spidermw.scrape_response_async( - self.call_spider, result, request, self.crawler.spider + self.call_spider_async, result, request, self.crawler.spider ) except Exception: self.handle_spider_error(Failure(), request, result) @@ -224,12 +224,11 @@ class Scraper: def call_spider( self, result: Response | Failure, request: Request, spider: Spider | None = None ) -> Deferred[Iterable[Any] | AsyncIterator[Any]]: - if spider is not None: - warnings.warn( - "Passing a 'spider' argument to Scraper.call_spider() is deprecated.", - category=ScrapyDeprecationWarning, - stacklevel=2, - ) + warnings.warn( + "Scraper.call_spider() is deprecated, use call_spider_async() instead", + ScrapyDeprecationWarning, + stacklevel=2, + ) return deferred_from_coro(self.call_spider_async(result, request)) async def call_spider_async( @@ -314,12 +313,11 @@ class Scraper: spider: Spider | None = None, ) -> Deferred[None]: """Pass items/requests produced by a callback to ``_process_spidermw_output()`` in parallel.""" - if spider is not None: - warnings.warn( - "Passing a 'spider' argument to Scraper.handle_spider_output() is deprecated.", - category=ScrapyDeprecationWarning, - stacklevel=2, - ) + warnings.warn( + "Scraper.handle_spider_output() is deprecated, use handle_spider_output_async() instead", + ScrapyDeprecationWarning, + stacklevel=2, + ) return deferred_from_coro( self.handle_spider_output_async(result, request, response) ) @@ -395,6 +393,11 @@ 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``. """ + warnings.warn( + "Scraper.start_itemproc() is deprecated, use start_itemproc_async() instead", + ScrapyDeprecationWarning, + stacklevel=2, + ) return deferred_from_coro(self.start_itemproc_async(item, response=response)) async def start_itemproc_async( diff --git a/scrapy/core/spidermw.py b/scrapy/core/spidermw.py index 01e563e56..a2513f4ec 100644 --- a/scrapy/core/spidermw.py +++ b/scrapy/core/spidermw.py @@ -7,7 +7,8 @@ See documentation in docs/topics/spider-middleware.rst from __future__ import annotations import logging -from collections.abc import AsyncIterator, Callable, Iterable +from collections.abc import AsyncIterator, Callable, Coroutine, Iterable +from functools import wraps from inspect import isasyncgenfunction, iscoroutine from itertools import islice from typing import TYPE_CHECKING, Any, TypeVar, Union, cast @@ -41,7 +42,7 @@ logger = logging.getLogger(__name__) _T = TypeVar("_T") ScrapeFunc = Callable[ [Union[Response, Failure], Request], - Deferred[Union[Iterable[_T], AsyncIterator[_T]]], + Coroutine[Any, Any, Union[Iterable[_T], AsyncIterator[_T]]], ] @@ -130,13 +131,13 @@ class SpiderMiddlewareManager(MiddlewareManager): process_spider_exception = getattr(mw, "process_spider_exception", None) self.methods["process_spider_exception"].appendleft(process_spider_exception) - def _process_spider_input( + async def _process_spider_input( self, scrape_func: ScrapeFunc[_T], response: Response, request: Request, spider: Spider, - ) -> Deferred[Iterable[_T] | AsyncIterator[_T]]: + ) -> Iterable[_T] | AsyncIterator[_T]: for method in self.methods["process_spider_input"]: method = cast("Callable", method) try: @@ -150,8 +151,8 @@ class SpiderMiddlewareManager(MiddlewareManager): except _InvalidOutput: raise except Exception: - return scrape_func(Failure(), request) - return scrape_func(response, request) + return await scrape_func(Failure(), request) + return await scrape_func(response, request) def _evaluate_iterable( self, @@ -362,13 +363,28 @@ class SpiderMiddlewareManager(MiddlewareManager): def scrape_response( self, - scrape_func: ScrapeFunc[_T], + scrape_func: Callable[ + [Response | Failure, Request], + Deferred[Iterable[_T] | AsyncIterator[_T]], + ], response: Response, request: Request, spider: Spider, ) -> Deferred[MutableChain[_T] | MutableAsyncChain[_T]]: + warn( + "SpiderMiddlewareManager.scrape_response() is deprecated, use scrape_response_async() instead", + ScrapyDeprecationWarning, + stacklevel=2, + ) + + @wraps(scrape_func) + async def scrape_func_wrapped( + response: Response | Failure, request: Request + ) -> Iterable[_T] | AsyncIterator[_T]: + return await maybe_deferred_to_future(scrape_func(response, request)) + return deferred_from_coro( - self.scrape_response_async(scrape_func, response, request, spider) + self.scrape_response_async(scrape_func_wrapped, response, request, spider) ) async def scrape_response_async( @@ -389,8 +405,8 @@ class SpiderMiddlewareManager(MiddlewareManager): return self._process_spider_exception(response, spider, exception) try: - it: Iterable[_T] | AsyncIterator[_T] = await maybe_deferred_to_future( - self._process_spider_input(scrape_func, response, request, spider) + it: Iterable[_T] | AsyncIterator[_T] = await self._process_spider_input( + scrape_func, response, request, spider ) return await process_callback_output(it) except Exception as ex: diff --git a/scrapy/crawler.py b/scrapy/crawler.py index c6c65a993..b8a2855f0 100644 --- a/scrapy/crawler.py +++ b/scrapy/crawler.py @@ -159,8 +159,8 @@ class Crawler: self._apply_settings() self._update_root_log_handler() self.engine = self._create_engine() - yield self.engine.open_spider(self.spider) - yield self.engine.start() + yield deferred_from_coro(self.engine.open_spider_async(self.spider)) + yield deferred_from_coro(self.engine.start_async()) except Exception: self.crawling = False if self.engine is not None: diff --git a/scrapy/shell.py b/scrapy/shell.py index c3a274e0d..311972c90 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, maybe_deferred_to_future +from scrapy.utils.defer import 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 @@ -126,9 +126,7 @@ class Shell: self.crawler.spider = spider assert self.crawler.engine - await maybe_deferred_to_future( - self.crawler.engine.open_spider(spider, close_if_idle=False) - ) + await self.crawler.engine.open_spider_async(spider, close_if_idle=False) self.crawler.engine._start_request_processing() self.spider = spider diff --git a/scrapy/utils/defer.py b/scrapy/utils/defer.py index fc149e185..5b65e68c0 100644 --- a/scrapy/utils/defer.py +++ b/scrapy/utils/defer.py @@ -511,3 +511,18 @@ def maybe_deferred_to_future(d: Deferred[_T]) -> Deferred[_T] | Future[_T]: if not is_asyncio_available(): return d return deferred_to_future(d) + + +def _schedule_coro(coro: Coroutine[Any, Any, Any]) -> None: + """Schedule the coroutine as a task or a Deferred. + + This doesn't store the reference to the task/Deferred, so a better + alternative is calling :func:`scrapy.utils.defer.deferred_from_coro`, + keeping the result, and adding proper exception handling (e.g. errbacks) to + it. + """ + if not is_asyncio_available(): + Deferred.fromCoroutine(coro) + return + loop = asyncio.get_event_loop() + loop.create_task(coro) # noqa: RUF006 diff --git a/tests/test_engine.py b/tests/test_engine.py index 14c8f8fee..a742fd365 100644 --- a/tests/test_engine.py +++ b/tests/test_engine.py @@ -10,6 +10,7 @@ module with the ``runserver`` argument:: python test_engine.py runserver """ +import asyncio import re import subprocess import sys @@ -36,8 +37,15 @@ from scrapy.item import Field, Item from scrapy.linkextractors import LinkExtractor from scrapy.signals import request_scheduled from scrapy.spiders import Spider -from scrapy.utils.defer import deferred_f_from_coro_f, maybe_deferred_to_future +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 +from scrapy.utils.spider import DefaultSpider from scrapy.utils.test import get_crawler from tests import get_testdata from tests.mockserver.http import MockServer @@ -433,13 +441,22 @@ class TestEngine(TestEngineBase): @inlineCallbacks def test_start_already_running_exception(self): - e = ExecutionEngine(get_crawler(MySpider), lambda _: None) - yield e.open_spider(MySpider()) - e.start() + e = ExecutionEngine(get_crawler(DefaultSpider), lambda _: None) + yield deferred_from_coro(e.open_spider_async(DefaultSpider())) + _schedule_coro(e.start_async()) with pytest.raises(RuntimeError, match="Engine already running"): - yield e.start() + yield deferred_from_coro(e.start_async()) yield e.stop() + @pytest.mark.only_asyncio + @deferred_f_from_coro_f + async def test_start_already_running_exception_asyncio(self): + e = ExecutionEngine(get_crawler(DefaultSpider), lambda _: None) + await e.open_spider_async(DefaultSpider()) + with pytest.raises(RuntimeError, match="Engine already running"): + await asyncio.gather(e.start_async(), e.start_async()) + await deferred_to_future(e.stop()) + @inlineCallbacks def test_start_request_processing_exception(self): class BadRequestFingerprinter: diff --git a/tests/test_spidermiddleware.py b/tests/test_spidermiddleware.py index bc2404d8b..f4811a741 100644 --- a/tests/test_spidermiddleware.py +++ b/tests/test_spidermiddleware.py @@ -31,7 +31,7 @@ class TestSpiderMiddleware: self.mwman = SpiderMiddlewareManager.from_crawler(self.crawler) async def _scrape_response(self) -> Any: - """Execute spider mw manager's scrape_response method and return the result. + """Execute spider mw manager's scrape_response_async method and return the result. Raise exception in case of failure. """ @@ -41,10 +41,8 @@ class TestSpiderMiddleware: it = mock.MagicMock() return defer.succeed(it) - return await maybe_deferred_to_future( - self.mwman.scrape_response( - scrape_func, self.response, self.request, self.spider - ) + return await self.mwman.scrape_response_async( + scrape_func, self.response, self.request, self.spider ) @@ -136,10 +134,10 @@ class TestBaseAsyncSpiderMiddleware(TestSpiderMiddleware): yield {"foo": 2} yield {"foo": 3} - def _scrape_func( + async def _scrape_func( self, response: Response | Failure, request: Request - ) -> defer.Deferred[Iterable[Any] | AsyncIterator[Any]]: - return defer.succeed(self._callback()) + ) -> Iterable[Any] | AsyncIterator[Any]: + return self._callback() async def _get_middleware_result( self, *mw_classes: type[Any], start_index: int | None = None