diff --git a/docs/topics/signals.rst b/docs/topics/signals.rst index 66cb87fc5..59742ffeb 100644 --- a/docs/topics/signals.rst +++ b/docs/topics/signals.rst @@ -46,8 +46,8 @@ Here is a simple example showing how you can catch signals and perform some acti .. _signal-deferred: -Deferred signal handlers -======================== +Asynchronous signal handlers +============================ Some signals support returning :class:`~twisted.internet.defer.Deferred` or :term:`awaitable objects ` from their handlers, allowing @@ -114,7 +114,7 @@ engine_started Sent when the Scrapy engine has started crawling. - This signal supports returning deferreds from its handlers. + This signal supports asynchronous handlers. .. note:: This signal may be fired *after* the :signal:`spider_opened` signal, depending on how the spider was started. So **don't** rely on this signal @@ -129,7 +129,7 @@ engine_stopped Sent when the Scrapy engine is stopped (for example, when a crawling process has finished). - This signal supports returning deferreds from its handlers. + This signal supports asynchronous handlers. scheduler_empty ~~~~~~~~~~~~~~~ @@ -164,7 +164,7 @@ item_scraped Sent when an item has been scraped, after it has passed all the :ref:`topics-item-pipeline` stages (without being dropped). - This signal supports returning deferreds from its handlers. + This signal supports asynchronous handlers. :param item: the scraped item :type item: :ref:`item object ` @@ -185,7 +185,7 @@ item_dropped Sent after an item has been dropped from the :ref:`topics-item-pipeline` when some stage raised a :exc:`~scrapy.exceptions.DropItem` exception. - This signal supports returning deferreds from its handlers. + This signal supports asynchronous handlers. :param item: the item dropped from the :ref:`topics-item-pipeline` :type item: :ref:`item object ` @@ -211,7 +211,7 @@ item_error Sent when a :ref:`topics-item-pipeline` generates an error (i.e. raises an exception), except :exc:`~scrapy.exceptions.DropItem` exception. - This signal supports returning deferreds from its handlers. + This signal supports asynchronous handlers. :param item: the item that caused the error in the :ref:`topics-item-pipeline` :type item: :ref:`item object ` @@ -239,7 +239,7 @@ spider_closed Sent after a spider has been closed. This can be used to release per-spider resources reserved on :signal:`spider_opened`. - This signal supports returning deferreds from its handlers. + This signal supports asynchronous handlers. :param spider: the spider which has been closed :type spider: :class:`~scrapy.Spider` object @@ -263,7 +263,7 @@ spider_opened reserve per-spider resources, but can be used for any task that needs to be performed when a spider is opened. - This signal supports returning deferreds from its handlers. + This signal supports asynchronous handlers. :param spider: the spider which has been opened :type spider: :class:`~scrapy.Spider` object @@ -332,7 +332,7 @@ feed_slot_closed Sent when a :ref:`feed exports ` slot is closed. - This signal supports returning deferreds from its handlers. + This signal supports asynchronous handlers. :param slot: the slot closed :type slot: scrapy.extensions.feedexport.FeedSlot @@ -348,7 +348,7 @@ feed_exporter_closed during the handling of the :signal:`spider_closed` signal by the extension, after all feed exporting has been handled. - This signal supports returning deferreds from its handlers. + This signal supports asynchronous handlers. Request signals diff --git a/scrapy/core/engine.py b/scrapy/core/engine.py index 658f6e774..b0d9a5452 100644 --- a/scrapy/core/engine.py +++ b/scrapy/core/engine.py @@ -127,9 +127,7 @@ class ExecutionEngine: if self.running: raise RuntimeError("Engine already running") self.start_time = time() - await maybe_deferred_to_future( - self.signals.send_catch_log_deferred(signal=signals.engine_started) - ) + await self.signals.send_catch_log_async(signal=signals.engine_started) self.running = True self._closewait: Deferred[None] = Deferred() if _start_request_processing: @@ -141,9 +139,7 @@ class ExecutionEngine: @deferred_f_from_coro_f async def _finish_stopping_engine(_: Any) -> None: - await maybe_deferred_to_future( - self.signals.send_catch_log_deferred(signal=signals.engine_stopped) - ) + await self.signals.send_catch_log_async(signal=signals.engine_stopped) self._closewait.callback(None) if not self.running: @@ -415,9 +411,7 @@ class ExecutionEngine: await maybe_deferred_to_future(self.scraper.open_spider(spider)) assert self.crawler.stats self.crawler.stats.open_spider(spider) - await maybe_deferred_to_future( - self.signals.send_catch_log_deferred(signals.spider_opened, spider=spider) - ) + await self.signals.send_catch_log_async(signals.spider_opened, spider=spider) def _spider_idle(self) -> None: """ diff --git a/scrapy/core/scraper.py b/scrapy/core/scraper.py index 2c48a9a81..9fc1d20ed 100644 --- a/scrapy/core/scraper.py +++ b/scrapy/core/scraper.py @@ -392,14 +392,12 @@ class Scraper: logger.log( *logformatter_adapter(logkws), extra={"spider": self.crawler.spider} ) - await maybe_deferred_to_future( - self.signals.send_catch_log_deferred( - signal=signals.item_dropped, - item=item, - response=response, - spider=self.crawler.spider, - exception=ex, - ) + await self.signals.send_catch_log_async( + signal=signals.item_dropped, + item=item, + response=response, + spider=self.crawler.spider, + exception=ex, ) except Exception as ex: logkws = self.logformatter.item_error( @@ -410,14 +408,12 @@ class Scraper: extra={"spider": self.crawler.spider}, exc_info=True, ) - await maybe_deferred_to_future( - self.signals.send_catch_log_deferred( - signal=signals.item_error, - item=item, - response=response, - spider=self.crawler.spider, - failure=Failure(), - ) + await self.signals.send_catch_log_async( + signal=signals.item_error, + item=item, + response=response, + spider=self.crawler.spider, + failure=Failure(), ) else: logkws = self.logformatter.scraped(output, response, self.crawler.spider) @@ -425,13 +421,11 @@ class Scraper: logger.log( *logformatter_adapter(logkws), extra={"spider": self.crawler.spider} ) - await maybe_deferred_to_future( - self.signals.send_catch_log_deferred( - signal=signals.item_scraped, - item=output, - response=response, - spider=self.crawler.spider, - ) + await self.signals.send_catch_log_async( + signal=signals.item_scraped, + item=output, + response=response, + spider=self.crawler.spider, ) finally: self.slot.itemproc_size -= 1 diff --git a/scrapy/extensions/feedexport.py b/scrapy/extensions/feedexport.py index 8bcd4e40d..c39a9c92e 100644 --- a/scrapy/extensions/feedexport.py +++ b/scrapy/extensions/feedexport.py @@ -531,9 +531,7 @@ class FeedExporter: await maybe_deferred_to_future(DeferredList(self._pending_deferreds)) # Send FEED_EXPORTER_CLOSED signal - await maybe_deferred_to_future( - self.crawler.signals.send_catch_log_deferred(signals.feed_exporter_closed) - ) + await self.crawler.signals.send_catch_log_async(signals.feed_exporter_closed) def _close_slot(self, slot: FeedSlot, spider: Spider) -> Deferred[None] | None: def get_file(slot_: FeedSlot) -> IO[bytes]: diff --git a/scrapy/signalmanager.py b/scrapy/signalmanager.py index f8c50b5e3..7fd172535 100644 --- a/scrapy/signalmanager.py +++ b/scrapy/signalmanager.py @@ -53,11 +53,10 @@ class SignalManager: self, signal: Any, **kwargs: Any ) -> Deferred[list[tuple[Any, Any]]]: """ - Like :meth:`send_catch_log` but supports returning - :class:`~twisted.internet.defer.Deferred` objects from signal handlers. + Like :meth:`send_catch_log` but supports asynchronous signal handlers. Returns a Deferred that gets fired once all signal handlers - deferreds were fired. Send a signal, catch exceptions and log them. + have finished. Send a signal, catch exceptions and log them. The keyword arguments are passed to the signal handlers (connected through the :meth:`connect` method). @@ -65,6 +64,21 @@ class SignalManager: kwargs.setdefault("sender", self.sender) return _signal.send_catch_log_deferred(signal, **kwargs) + async def send_catch_log_async( + self, signal: Any, **kwargs: Any + ) -> list[tuple[Any, Any]]: + """ + Like :meth:`send_catch_log` but supports asynchronous signal handlers. + + Returns a coroutine that completes once all signal handlers + have finished. Send a signal, catch exceptions and log them. + + The keyword arguments are passed to the signal handlers (connected + through the :meth:`connect` method). + """ + kwargs.setdefault("sender", self.sender) + return await _signal.send_catch_log_async(signal, **kwargs) + def disconnect_all(self, signal: Any, **kwargs: Any) -> None: """ Disconnect all receivers from the given signal. diff --git a/scrapy/utils/signal.py b/scrapy/utils/signal.py index 5fd176a3f..d6b0a671b 100644 --- a/scrapy/utils/signal.py +++ b/scrapy/utils/signal.py @@ -3,7 +3,7 @@ from __future__ import annotations import logging -from collections.abc import Sequence +from collections.abc import Generator, Sequence from typing import Any as TypingAny from pydispatch.dispatcher import ( @@ -14,11 +14,11 @@ from pydispatch.dispatcher import ( liveReceivers, ) from pydispatch.robustapply import robustApply -from twisted.internet.defer import Deferred, DeferredList +from twisted.internet.defer import Deferred, DeferredList, inlineCallbacks from twisted.python.failure import Failure from scrapy.exceptions import StopDownload -from scrapy.utils.defer import maybeDeferred_coro +from scrapy.utils.defer import maybe_deferred_to_future, maybeDeferred_coro from scrapy.utils.log import failure_to_exc_info logger = logging.getLogger(__name__) @@ -66,18 +66,19 @@ def send_catch_log( return responses +@inlineCallbacks def send_catch_log_deferred( signal: TypingAny = Any, sender: TypingAny = Anonymous, *arguments: TypingAny, **named: TypingAny, -) -> Deferred[list[tuple[TypingAny, TypingAny]]]: - """Like send_catch_log but supports returning deferreds on signal handlers. - Returns a deferred that gets fired once all signal handlers deferreds were - fired. +) -> Generator[Deferred[TypingAny], TypingAny, list[tuple[TypingAny, TypingAny]]]: + """Like send_catch_log but supports asynchronous signal handlers. + + Returns a deferred that gets fired once all signal handlers have finished. """ - def logerror(failure: Failure, recv: Any) -> Failure: + def logerror(failure: Failure, recv: TypingAny) -> Failure: if dont_log is None or not isinstance(failure.value, dont_log): logger.error( "Error caught on signal handler: %(receiver)s", @@ -103,11 +104,24 @@ def send_catch_log_deferred( ) ) dfds.append(d2) - dl = DeferredList(dfds) - d3: Deferred[list[tuple[TypingAny, TypingAny]]] = dl.addCallback( - lambda out: [x[1] for x in out] + + results = yield DeferredList(dfds) + return [result[1] for result in results] + + +async def send_catch_log_async( + signal: TypingAny = Any, + sender: TypingAny = Anonymous, + *arguments: TypingAny, + **named: TypingAny, +) -> list[tuple[TypingAny, TypingAny]]: + """Like send_catch_log but supports asynchronous signal handlers. + + Returns a coroutine that completes once all signal handlers have finished. + """ + return await maybe_deferred_to_future( + send_catch_log_deferred(signal, sender, *arguments, **named) ) - return d3 def disconnect_all(signal: TypingAny = Any, sender: TypingAny = Any) -> None: diff --git a/tests/test_utils_signal.py b/tests/test_utils_signal.py index 751a77031..6dff321da 100644 --- a/tests/test_utils_signal.py +++ b/tests/test_utils_signal.py @@ -7,7 +7,12 @@ from twisted.internet import defer, reactor from twisted.python.failure import Failure from twisted.trial import unittest -from scrapy.utils.signal import send_catch_log, send_catch_log_deferred +from scrapy.utils.defer import deferred_from_coro +from scrapy.utils.signal import ( + send_catch_log, + send_catch_log_async, + send_catch_log_deferred, +) from scrapy.utils.test import get_from_asyncio_queue @@ -85,6 +90,38 @@ class SendCatchLogDeferredAsyncioTest(SendCatchLogDeferredTest): return await get_from_asyncio_queue("OK") +class SendCatchLogAsyncTest(TestSendCatchLog): + def _get_result(self, signal, *a, **kw): + return deferred_from_coro(send_catch_log_async(signal, *a, **kw)) + + +class SendCatchLogAsyncTest2(SendCatchLogAsyncTest): + def ok_handler(self, arg, handlers_called): + handlers_called.add(self.ok_handler) + assert arg == "test" + d = defer.Deferred() + reactor.callLater(0, d.callback, "OK") + return d + + +@pytest.mark.usefixtures("reactor_pytest") +class SendCatchLogAsyncAsyncDefTest(SendCatchLogAsyncTest): + async def ok_handler(self, arg, handlers_called): + handlers_called.add(self.ok_handler) + assert arg == "test" + await defer.succeed(42) + return "OK" + + +@pytest.mark.only_asyncio +class SendCatchLogAsyncAsyncioTest(SendCatchLogAsyncTest): + async def ok_handler(self, arg, handlers_called): + handlers_called.add(self.ok_handler) + assert arg == "test" + await asyncio.sleep(0.2) + return await get_from_asyncio_queue("OK") + + class TestSendCatchLog2: def test_error_logged_if_deferred_not_supported(self): def test_handler():