Add send_catch_log_async().

This commit is contained in:
Andrey Rakhmatullin 2025-05-15 14:18:01 +05:00
parent 82acef3051
commit 1ddcb568e2
7 changed files with 113 additions and 62 deletions

View File

@ -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 <awaitable>` 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 <item-types>`
@ -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 <item-types>`
@ -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 <item-types>`
@ -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 <topics-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

View File

@ -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:
"""

View File

@ -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

View File

@ -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]:

View File

@ -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.

View File

@ -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:

View File

@ -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():