Deprecate send_catch_log_deferred(). (#7161)

* Add a test for not having pending tasks.

* Refactor TestFeedExporterSignals.

* Refactor FeedExporter.close_spider().

* More engine start/stop robustness.

* Refactor send_catch_log_async(), deprecate send_catch_log_deferred().

* Add pragma: no cover.

* Warn on signal handlers returning a Deferred.

* Make _pending_close_coros an instance attribute.

* Remove an unused function.

* Remove the outdated TODO.
This commit is contained in:
Andrey Rakhmatullin 2025-12-10 14:42:49 +05:00 committed by GitHub
parent 5105f55a98
commit 9bfa58e36c
No known key found for this signature in database
GPG Key ID: B5690EEEBB952194
10 changed files with 204 additions and 103 deletions

View File

@ -74,7 +74,7 @@ In addition to native coroutine APIs Scrapy has some APIs that return a
function that returns a :class:`~twisted.internet.defer.Deferred` object. These
APIs are also asynchronous but don't yet support native ``async def`` syntax.
In the future we plan to add support for the ``async def`` syntax to these APIs
or replace them with other APIs where changing the existing ones is
or replace them with other APIs where changing the existing ones isn't
possible.
These APIs don't have a coroutine-based counterpart:
@ -99,14 +99,6 @@ These APIs have a coroutine-based implementation and a Deferred-based one:
doesn't support non-default reactors and so the latter should be used
with those.
- :class:`scrapy.signalmanager.SignalManager`:
- :meth:`~scrapy.signalmanager.SignalManager.send_catch_log_async`
(coroutine-based) and
:meth:`~scrapy.signalmanager.SignalManager.send_catch_log_deferred`
(Deferred-based): the latter will be deprecated in a later Scrapy
version.
The following user-supplied methods can return
:class:`~twisted.internet.defer.Deferred` objects (the methods that can also
return coroutines are listed in :ref:`coroutine-support`):

View File

@ -5,17 +5,16 @@ from __future__ import annotations
import logging
from typing import TYPE_CHECKING, Any, Protocol, cast
from twisted.internet import defer
from scrapy import Request, Spider, signals
from scrapy.exceptions import NotConfigured, NotSupported
from scrapy.utils.decorators import _warn_spider_arg
from scrapy.utils.defer import ensure_awaitable
from scrapy.utils.httpobj import urlparse_cached
from scrapy.utils.misc import build_from_crawler, load_object
from scrapy.utils.python import without_none_values
if TYPE_CHECKING:
from collections.abc import Callable, Generator
from collections.abc import Callable
from twisted.internet.defer import Deferred
@ -107,8 +106,9 @@ class DownloadHandlers:
assert self._crawler.spider
return handler.download_request(request, self._crawler.spider)
@defer.inlineCallbacks
def _close(self, *_a: Any, **_kw: Any) -> Generator[Deferred[Any], Any, None]:
async def _close(self) -> None:
for dh in self._handlers.values():
if hasattr(dh, "close"):
yield dh.close()
if not hasattr(dh, "close"):
continue
await ensure_awaitable(dh.close())

View File

@ -115,6 +115,8 @@ class ExecutionEngine:
self._slot: _Slot | None = None
self.spider: Spider | None = None
self.running: bool = False
self._starting: bool = False
self._stopping: bool = False
self.paused: bool = False
self._spider_closed_callback: Callable[
[Spider], Coroutine[Any, Any, None] | Deferred[None] | None
@ -172,10 +174,14 @@ class ExecutionEngine:
.. versionadded:: VERSION
"""
if self.running:
if self._starting:
raise RuntimeError("Engine already running")
self.start_time = time()
self._starting = True
await self.signals.send_catch_log_async(signal=signals.engine_started)
if self._stopping:
# band-aid until https://github.com/scrapy/scrapy/issues/6916
return
if _start_request_processing and self.spider is None:
# require an opened spider when not run in scrapy shell
return
@ -205,12 +211,21 @@ class ExecutionEngine:
.. versionadded:: VERSION
"""
if not self.running:
if not self._starting:
raise RuntimeError("Engine not running")
self.running = False
self.running = self._starting = False
self._stopping = True
if self._start_request_processing_awaitable is not None:
self._start_request_processing_awaitable.cancel()
if (
not is_asyncio_available()
or self._start_request_processing_awaitable
is not asyncio.current_task()
):
# If using the asyncio loop and stop_async() was called from
# start() itself, we can't cancel it, and _start_request_processing()
# will exit via the self.running check.
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")
@ -285,7 +300,7 @@ class ExecutionEngine:
self._slot.nextcall.schedule()
self._slot.heartbeat.start(self._SLOT_HEARTBEAT_INTERVAL)
while self._start and self.spider:
while self._start and self.spider and self.running:
await self._process_start_next()
if not self.needs_backout():
# Give room for the outcome of self._process_start_next() to be
@ -293,7 +308,7 @@ class ExecutionEngine:
self._slot.nextcall.schedule()
await self._slot.nextcall.wait()
except (asyncio.exceptions.CancelledError, CancelledError):
# self.stop() has cancelled us, nothing to do
# self.stop_async() has cancelled us, nothing to do
return
except Exception:
# an error happened, log it and stop the engine

View File

@ -6,20 +6,21 @@ See documentation in docs/topics/feed-exports.rst
from __future__ import annotations
import asyncio
import contextlib
import logging
import re
import sys
import warnings
from abc import ABC, abstractmethod
from collections.abc import Callable
from collections.abc import Callable, Coroutine
from datetime import datetime, timezone
from pathlib import Path, PureWindowsPath
from tempfile import NamedTemporaryFile
from typing import IO, TYPE_CHECKING, Any, Protocol, TypeAlias, cast
from urllib.parse import unquote, urlparse
from twisted.internet.defer import Deferred, DeferredList, maybeDeferred
from twisted.internet.defer import Deferred, DeferredList
from twisted.internet.threads import deferToThread
from w3lib.url import file_uri_to_path
from zope.interface import Interface, implementer
@ -27,16 +28,15 @@ from zope.interface import Interface, implementer
from scrapy import Spider, signals
from scrapy.exceptions import NotConfigured, ScrapyDeprecationWarning
from scrapy.extensions.postprocessing import PostProcessingManager
from scrapy.utils.asyncio import is_asyncio_available
from scrapy.utils.conf import feed_complete_default_values_from_settings
from scrapy.utils.defer import maybe_deferred_to_future
from scrapy.utils.defer import deferred_from_coro, ensure_awaitable
from scrapy.utils.ftp import ftp_store_file
from scrapy.utils.log import failure_to_exc_info
from scrapy.utils.misc import build_from_crawler, load_object
from scrapy.utils.python import without_none_values
if TYPE_CHECKING:
from _typeshed import OpenBinaryMode
from twisted.python.failure import Failure
# typing.Self requires Python 3.11
from typing_extensions import Self
@ -431,8 +431,6 @@ class FeedSlot:
class FeedExporter:
_pending_deferreds: list[Deferred[None]] = []
@classmethod
def from_crawler(cls, crawler: Crawler) -> Self:
exporter = cls(crawler)
@ -447,6 +445,7 @@ class FeedExporter:
self.feeds = {}
self.slots: list[FeedSlot] = []
self.filters: dict[str, ItemFilter] = {}
self._pending_close_coros: list[Coroutine[Any, Any, None]] = []
if not self.settings["FEEDS"] and not self.settings["FEED_URI"]:
raise NotConfigured
@ -506,17 +505,24 @@ class FeedExporter:
)
async def close_spider(self, spider: Spider) -> None:
for slot in self.slots:
self._close_slot(slot, spider)
self._pending_close_coros.extend(
self._close_slot(slot, spider) for slot in self.slots
)
# Await all deferreds
if self._pending_deferreds:
await maybe_deferred_to_future(DeferredList(self._pending_deferreds))
if self._pending_close_coros:
if is_asyncio_available():
await asyncio.wait(
[asyncio.create_task(coro) for coro in self._pending_close_coros]
)
else:
await DeferredList(
deferred_from_coro(coro) for coro in self._pending_close_coros
)
# Send FEED_EXPORTER_CLOSED signal
await self.crawler.signals.send_catch_log_async(signals.feed_exporter_closed)
def _close_slot(self, slot: FeedSlot, spider: Spider) -> Deferred[None] | None:
async def _close_slot(self, slot: FeedSlot, spider: Spider) -> None:
def get_file(slot_: FeedSlot) -> IO[bytes]:
assert slot_.file
if isinstance(slot_.file, PostProcessingManager):
@ -533,45 +539,28 @@ class FeedExporter:
slot.finish_exporting()
else:
# In this case, the file is not stored, so no processing is required.
return None
return
logmsg = f"{slot.format} feed ({slot.itemcount} items) in: {slot.uri}"
d: Deferred[None] = maybeDeferred(slot.storage.store, get_file(slot)) # type: ignore[call-overload]
d.addCallback(
self._handle_store_success, logmsg, spider, type(slot.storage).__name__
)
d.addErrback(
self._handle_store_error, logmsg, spider, type(slot.storage).__name__
)
self._pending_deferreds.append(d)
d.addCallback(
lambda _: self.crawler.signals.send_catch_log_deferred(
signals.feed_slot_closed, slot=slot
slot_type = type(slot.storage).__name__
assert self.crawler.stats
try:
await ensure_awaitable(slot.storage.store(get_file(slot)))
except Exception:
logger.error(
"Error storing %s",
logmsg,
exc_info=True,
extra={"spider": spider},
)
self.crawler.stats.inc_value(f"feedexport/failed_count/{slot_type}")
else:
logger.info("Stored %s", logmsg, extra={"spider": spider})
self.crawler.stats.inc_value(f"feedexport/success_count/{slot_type}")
await self.crawler.signals.send_catch_log_async(
signals.feed_slot_closed, slot=slot
)
d.addBoth(lambda _: self._pending_deferreds.remove(d))
return d
def _handle_store_error(
self, f: Failure, logmsg: str, spider: Spider, slot_type: str
) -> None:
logger.error(
"Error storing %s",
logmsg,
exc_info=failure_to_exc_info(f),
extra={"spider": spider},
)
assert self.crawler.stats
self.crawler.stats.inc_value(f"feedexport/failed_count/{slot_type}")
def _handle_store_success(
self, result: Any, logmsg: str, spider: Spider, slot_type: str
) -> None:
logger.info("Stored %s", logmsg, extra={"spider": spider})
assert self.crawler.stats
self.crawler.stats.inc_value(f"feedexport/success_count/{slot_type}")
def _start_new_batch(
self,
@ -627,7 +616,7 @@ class FeedExporter:
uri_params = self._get_uri_params(
spider, self.feeds[slot.uri_template]["uri_params"], slot
)
self._close_slot(slot, spider)
self._pending_close_coros.append(self._close_slot(slot, spider))
slots.append(
self._start_new_batch(
batch_id=slot.batch_id + 1,

View File

@ -1,10 +1,12 @@
from __future__ import annotations
import warnings
from typing import Any
from pydispatch import dispatcher
from twisted.internet.defer import Deferred
from scrapy.exceptions import ScrapyDeprecationWarning
from scrapy.utils import signal as _signal
from scrapy.utils.defer import maybe_deferred_to_future
@ -51,7 +53,7 @@ class SignalManager:
def send_catch_log_deferred(
self, signal: Any, **kwargs: Any
) -> Deferred[list[tuple[Any, Any]]]:
) -> Deferred[list[tuple[Any, Any]]]: # pragma: no cover
"""
Like :meth:`send_catch_log` but supports :ref:`asynchronous signal
handlers <signal-deferred>`.
@ -63,7 +65,12 @@ class SignalManager:
through the :meth:`connect` method).
"""
kwargs.setdefault("sender", self.sender)
return _signal.send_catch_log_deferred(signal, **kwargs)
warnings.warn(
"send_catch_log_deferred() is deprecated, use send_catch_log_async() instead",
ScrapyDeprecationWarning,
stacklevel=2,
)
return _signal._send_catch_log_deferred(signal, **kwargs)
async def send_catch_log_async(
self, signal: Any, **kwargs: Any
@ -80,6 +87,7 @@ class SignalManager:
.. versionadded:: VERSION
"""
# note that this returns exceptions instead of Failures in the second tuple member
kwargs.setdefault("sender", self.sender)
return await _signal.send_catch_log_async(signal, **kwargs)

View File

@ -2,8 +2,10 @@
from __future__ import annotations
import asyncio
import logging
from collections.abc import Generator, Sequence
import warnings
from collections.abc import Awaitable, Callable, Generator, Sequence
from typing import Any as TypingAny
from pydispatch.dispatcher import (
@ -17,9 +19,15 @@ from pydispatch.robustapply import robustApply
from twisted.internet.defer import Deferred, DeferredList, inlineCallbacks
from twisted.python.failure import Failure
from scrapy.exceptions import StopDownload
from scrapy.utils.defer import maybe_deferred_to_future, maybeDeferred_coro
from scrapy.exceptions import ScrapyDeprecationWarning, StopDownload
from scrapy.utils.asyncio import is_asyncio_available
from scrapy.utils.defer import (
ensure_awaitable,
maybe_deferred_to_future,
maybeDeferred_coro,
)
from scrapy.utils.log import failure_to_exc_info
from scrapy.utils.python import global_object_name
logger = logging.getLogger(__name__)
@ -66,19 +74,32 @@ def send_catch_log(
return responses
@inlineCallbacks
def send_catch_log_deferred(
signal: TypingAny = Any,
sender: TypingAny = Anonymous,
*arguments: TypingAny,
**named: TypingAny,
) -> Generator[Deferred[TypingAny], TypingAny, list[tuple[TypingAny, TypingAny]]]:
) -> Deferred[list[tuple[TypingAny, TypingAny]]]:
"""Like :func:`send_catch_log` but supports :ref:`asynchronous signal handlers
<signal-deferred>`.
Returns a deferred that gets fired once all signal handlers have finished.
"""
warnings.warn(
"send_catch_log_deferred() is deprecated, use send_catch_log_async() instead",
ScrapyDeprecationWarning,
stacklevel=2,
)
return _send_catch_log_deferred(signal, sender, *arguments, **named)
@inlineCallbacks
def _send_catch_log_deferred(
signal: TypingAny,
sender: TypingAny,
*arguments: TypingAny,
**named: TypingAny,
) -> Generator[Deferred[TypingAny], TypingAny, list[tuple[TypingAny, TypingAny]]]:
def logerror(failure: Failure, recv: TypingAny) -> Failure:
if dont_log is None or not isinstance(failure.value, dont_log):
logger.error(
@ -123,9 +144,65 @@ async def send_catch_log_async(
.. versionadded:: VERSION
"""
return await maybe_deferred_to_future(
send_catch_log_deferred(signal, sender, *arguments, **named)
# note that this returns exceptions instead of Failures in the second tuple member
if is_asyncio_available():
return await _send_catch_log_asyncio(signal, sender, *arguments, **named)
results = await maybe_deferred_to_future(
_send_catch_log_deferred(signal, sender, *arguments, **named)
)
return [
(receiver, result.value if isinstance(result, Failure) else result)
for receiver, result in results
]
async def _send_catch_log_asyncio(
signal: TypingAny = Any,
sender: TypingAny = Anonymous,
*arguments: TypingAny,
**named: TypingAny,
) -> list[tuple[TypingAny, TypingAny]]:
"""Like :func:`send_catch_log` but supports :ref:`asynchronous signal handlers
<signal-deferred>`.
Returns a coroutine that completes once all signal handlers have finished.
This function requires
:class:`~twisted.internet.asyncioreactor.AsyncioSelectorReactor` to be
installed.
.. versionadded:: VERSION
"""
dont_log = named.pop("dont_log", ())
dont_log = tuple(dont_log) if isinstance(dont_log, Sequence) else (dont_log,)
spider = named.get("spider")
handlers: list[Awaitable[TypingAny]] = []
for receiver in liveReceivers(getAllReceivers(sender, signal)):
async def handler(receiver: Callable) -> TypingAny:
result: TypingAny
try:
result = await ensure_awaitable(
robustApply(
receiver, signal=signal, sender=sender, *arguments, **named
),
_warn=global_object_name(receiver),
)
except dont_log as ex: # pylint: disable=catching-non-exception
result = ex
except Exception as ex:
logger.error(
"Error caught on signal handler: %(receiver)s",
{"receiver": receiver},
exc_info=True,
extra={"spider": spider},
)
result = ex
return (receiver, result)
handlers.append(handler(receiver))
return await asyncio.gather(*handlers, return_exceptions=True)
def disconnect_all(signal: TypingAny = Any, sender: TypingAny = Any) -> None:

View File

@ -659,6 +659,7 @@ class TestEngineCloseSpider:
engine = ExecutionEngine(crawler, lambda _: None)
with pytest.raises(RuntimeError, match="Spider not opened"):
await engine.close_spider_async()
engine.downloader.close() # cleanup
@deferred_f_from_coro_f
async def test_exception_slot(

View File

@ -29,7 +29,6 @@ import lxml.etree
import pytest
from packaging.version import Version
from testfixtures import LogCapture
from twisted.internet import defer
from twisted.internet.defer import inlineCallbacks
from w3lib.url import file_uri_to_path, path_to_file_uri
from zope.interface import implementer
@ -59,7 +58,7 @@ from tests.mockserver.http import MockServer
from tests.spiders import ItemSpider
if TYPE_CHECKING:
from collections.abc import Iterable
from collections.abc import Callable, Iterable
from os import PathLike
@ -2762,6 +2761,12 @@ class TestBatchDeliveries(TestFeedExportBase):
assert len(CustomS3FeedStorage.stubs) == len(items)
for stub in CustomS3FeedStorage.stubs[:-1]:
stub.assert_no_pending_responses()
assert (
"feedexport/success_count/CustomS3FeedStorage" in crawler.stats.get_stats()
)
assert (
crawler.stats.get_value("feedexport/success_count/CustomS3FeedStorage") == 3
)
# Test that the FeedExporer sends the feed_exporter_closed and feed_slot_closed signals
@ -2787,21 +2792,15 @@ class TestFeedExporterSignals:
def feed_slot_closed_signal_handler(self, slot):
self.feed_slot_closed_received = True
def feed_exporter_closed_signal_handler_deferred(self):
d = defer.Deferred()
d.addCallback(lambda _: setattr(self, "feed_exporter_closed_received", True))
d.callback(None)
return d
async def feed_exporter_closed_signal_handler_async(self):
self.feed_exporter_closed_received = True
def feed_slot_closed_signal_handler_deferred(self, slot):
d = defer.Deferred()
d.addCallback(lambda _: setattr(self, "feed_slot_closed_received", True))
d.callback(None)
return d
async def feed_slot_closed_signal_handler_async(self, slot):
self.feed_slot_closed_received = True
def run_signaled_feed_exporter(
self, feed_exporter_signal_handler, feed_slot_signal_handler
):
async def run_signaled_feed_exporter(
self, feed_exporter_signal_handler: Callable, feed_slot_signal_handler: Callable
) -> None:
crawler = get_crawler(settings_dict=self.settings)
feed_exporter = FeedExporter.from_crawler(crawler)
spider = scrapy.Spider("default")
@ -2816,26 +2815,28 @@ class TestFeedExporterSignals:
feed_exporter.open_spider(spider)
for item in self.items:
feed_exporter.item_scraped(item, spider)
defer.ensureDeferred(feed_exporter.close_spider(spider))
await feed_exporter.close_spider(spider)
def test_feed_exporter_signals_sent(self):
@deferred_f_from_coro_f
async def test_feed_exporter_signals_sent(self) -> None:
self.feed_exporter_closed_received = False
self.feed_slot_closed_received = False
self.run_signaled_feed_exporter(
await self.run_signaled_feed_exporter(
self.feed_exporter_closed_signal_handler,
self.feed_slot_closed_signal_handler,
)
assert self.feed_slot_closed_received
assert self.feed_exporter_closed_received
def test_feed_exporter_signals_sent_deferred(self):
@deferred_f_from_coro_f
async def test_feed_exporter_signals_sent_async(self) -> None:
self.feed_exporter_closed_received = False
self.feed_slot_closed_received = False
self.run_signaled_feed_exporter(
self.feed_exporter_closed_signal_handler_deferred,
self.feed_slot_closed_signal_handler_deferred,
await self.run_signaled_feed_exporter(
self.feed_exporter_closed_signal_handler_async,
self.feed_slot_closed_signal_handler_async,
)
assert self.feed_slot_closed_received
assert self.feed_exporter_closed_received

View File

@ -18,6 +18,9 @@ from scrapy.utils.test import get_from_asyncio_queue
class TestSendCatchLog:
# whether the function being tested returns exceptions or failures
returns_exceptions: bool = False
@inlineCallbacks
def test_send_catch_log(self):
test_signal = object()
@ -40,7 +43,9 @@ class TestSendCatchLog:
assert "error_handler" in record.getMessage()
assert record.levelname == "ERROR"
assert result[0][0] == self.error_handler # pylint: disable=comparison-with-callable
assert isinstance(result[0][1], Failure)
assert isinstance(
result[0][1], Exception if self.returns_exceptions else Failure
)
assert result[1] == (self.ok_handler, "OK")
dispatcher.disconnect(self.error_handler, signal=test_signal)
@ -59,6 +64,7 @@ class TestSendCatchLog:
return "OK"
@pytest.mark.filterwarnings("ignore::scrapy.exceptions.ScrapyDeprecationWarning")
class TestSendCatchLogDeferred(TestSendCatchLog):
def _get_result(self, signal, *a, **kw):
return send_catch_log_deferred(signal, *a, **kw)
@ -91,10 +97,13 @@ class TestSendCatchLogDeferredAsyncio(TestSendCatchLogDeferred):
class TestSendCatchLogAsync(TestSendCatchLog):
returns_exceptions = True
def _get_result(self, signal, *a, **kw):
return deferred_from_coro(send_catch_log_async(signal, *a, **kw))
@pytest.mark.filterwarnings("ignore::scrapy.exceptions.ScrapyDeprecationWarning")
class TestSendCatchLogAsync2(TestSendCatchLogAsync):
def ok_handler(self, arg, handlers_called):
handlers_called.add(self.ok_handler)

View File

@ -2,8 +2,11 @@
from __future__ import annotations
import asyncio
import logging
import pytest
from scrapy.utils.log import LogCounterHandler
@ -25,3 +28,9 @@ def test_stderr_log_handler() -> None:
"""
c = sum(1 for h in logging.root.handlers if type(h) is logging.StreamHandler) # pylint: disable=unidiomatic-typecheck
assert c == 0
@pytest.mark.only_asyncio
def test_pending_asyncio_tasks() -> None:
"""Test that there are no pending asyncio tasks."""
assert not asyncio.all_tasks()