From b50c032ee9a75d1c9b42f1126637fdc655b141a8 Mon Sep 17 00:00:00 2001 From: guillermo-bondonno <95530227+guillermo-bondonno@users.noreply.github.com> Date: Wed, 26 Apr 2023 03:20:37 -0300 Subject: [PATCH] Add feed_slot_closed and feed_exporter_closed signals (#5876) --- docs/topics/signals.rst | 27 ++++++++++ scrapy/extensions/feedexport.py | 43 +++++++++++---- scrapy/signals.py | 2 + tests/test_feedexport.py | 92 ++++++++++++++++++++++++++++++--- 4 files changed, 148 insertions(+), 16 deletions(-) diff --git a/docs/topics/signals.rst b/docs/topics/signals.rst index 3400a205a..9bfd1761c 100644 --- a/docs/topics/signals.rst +++ b/docs/topics/signals.rst @@ -307,6 +307,33 @@ spider_error :param spider: the spider which raised the exception :type spider: :class:`~scrapy.Spider` object +feed_slot_closed +~~~~~~~~~~~~~~~~ + +.. signal:: feed_slot_closed +.. function:: feed_slot_closed(slot) + + Sent when a :ref:`feed exports ` slot is closed. + + This signal supports returning deferreds from its handlers. + + :param slot: the slot closed + :type slot: scrapy.extensions.feedexport.FeedSlot + + +feed_exporter_closed +~~~~~~~~~~~~~~~~~~~~ + +.. signal:: feed_exporter_closed +.. function:: feed_exporter_closed() + + Sent when the :ref:`feed exports ` extension is 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. + + Request signals --------------- diff --git a/scrapy/extensions/feedexport.py b/scrapy/extensions/feedexport.py index da1a88299..bcf0b779a 100644 --- a/scrapy/extensions/feedexport.py +++ b/scrapy/extensions/feedexport.py @@ -11,10 +11,11 @@ import warnings from datetime import datetime from pathlib import Path from tempfile import NamedTemporaryFile -from typing import IO, Any, Callable, Optional, Tuple, Union +from typing import IO, Any, Callable, List, Optional, Tuple, Union from urllib.parse import unquote, urlparse from twisted.internet import defer, threads +from twisted.internet.defer import DeferredList from w3lib.url import file_uri_to_path from zope.interface import Interface, implementer @@ -23,6 +24,8 @@ from scrapy.exceptions import NotConfigured, ScrapyDeprecationWarning from scrapy.extensions.postprocessing import PostProcessingManager from scrapy.utils.boto import is_botocore_available from scrapy.utils.conf import feed_complete_default_values_from_settings +from scrapy.utils.defer import maybe_deferred_to_future +from scrapy.utils.deprecate import create_deprecated_class from scrapy.utils.ftp import ftp_store_file from scrapy.utils.log import failure_to_exc_info from scrapy.utils.misc import create_instance, load_object @@ -271,7 +274,7 @@ class FTPFeedStorage(BlockingFeedStorage): ) -class _FeedSlot: +class FeedSlot: def __init__( self, file, @@ -309,7 +312,15 @@ class _FeedSlot: self._exporting = False +_FeedSlot = create_deprecated_class( + name="_FeedSlot", + new_class=FeedSlot, +) + + class FeedExporter: + _pending_deferreds: List[defer.Deferred] = [] + @classmethod def from_crawler(cls, crawler): exporter = cls(crawler) @@ -375,12 +386,18 @@ class FeedExporter: ) ) - def close_spider(self, spider): - deferred_list = [] + async def close_spider(self, spider): for slot in self.slots: - d = self._close_slot(slot, spider) - deferred_list.append(d) - return defer.DeferredList(deferred_list) if deferred_list else None + self._close_slot(slot, spider) + + # Await all deferreds + if self._pending_deferreds: + 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) + ) def _close_slot(self, slot, spider): def get_file(slot_): @@ -404,6 +421,14 @@ class FeedExporter: 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 + ) + ) + d.addBoth(lambda _: self._pending_deferreds.remove(d)) + return d def _handle_store_error(self, f, logmsg, spider, slot_type): @@ -444,7 +469,7 @@ class FeedExporter: indent=feed_options["indent"], **feed_options["item_export_kwargs"], ) - slot = _FeedSlot( + slot = FeedSlot( file=file, exporter=exporter, storage=storage, @@ -579,7 +604,7 @@ class FeedExporter: self, spider: Spider, uri_params_function: Optional[Union[str, Callable[[dict, Spider], dict]]], - slot: Optional[_FeedSlot] = None, + slot: Optional[FeedSlot] = None, ) -> dict: params = {} for k in dir(spider): diff --git a/scrapy/signals.py b/scrapy/signals.py index 8cf2a4d93..0090f1c8b 100644 --- a/scrapy/signals.py +++ b/scrapy/signals.py @@ -22,6 +22,8 @@ bytes_received = object() item_scraped = object() item_dropped = object() item_error = object() +feed_slot_closed = object() +feed_exporter_closed = object() # for backward compatibility stats_spider_opened = spider_opened diff --git a/tests/test_feedexport.py b/tests/test_feedexport.py index 83de0e77e..b1059099a 100644 --- a/tests/test_feedexport.py +++ b/tests/test_feedexport.py @@ -32,18 +32,19 @@ from zope.interface import implementer from zope.interface.verify import verifyObject import scrapy +from scrapy import signals from scrapy.exceptions import NotConfigured, ScrapyDeprecationWarning from scrapy.exporters import CsvItemExporter, JsonItemExporter from scrapy.extensions.feedexport import ( BlockingFeedStorage, FeedExporter, + FeedSlot, FileFeedStorage, FTPFeedStorage, GCSFeedStorage, IFeedStorage, S3FeedStorage, StdoutFeedStorage, - _FeedSlot, ) from scrapy.settings import Settings from scrapy.utils.python import to_unicode @@ -660,8 +661,8 @@ class FeedExportTestBase(ABC, unittest.TestCase): return result -class InstrumentedFeedSlot(_FeedSlot): - """Instrumented _FeedSlot subclass for keeping track of calls to +class InstrumentedFeedSlot(FeedSlot): + """Instrumented FeedSlot subclass for keeping track of calls to start_exporting and finish_exporting.""" def start_exporting(self): @@ -964,7 +965,7 @@ class FeedExportTest(FeedExportTestBase): listener = IsExportingListener() InstrumentedFeedSlot.subscribe__listener(listener) - with mock.patch("scrapy.extensions.feedexport._FeedSlot", InstrumentedFeedSlot): + with mock.patch("scrapy.extensions.feedexport.FeedSlot", InstrumentedFeedSlot): _ = yield self.exported_data(items, settings) self.assertFalse(listener.start_without_finish) self.assertFalse(listener.finish_without_start) @@ -982,7 +983,7 @@ class FeedExportTest(FeedExportTestBase): listener = IsExportingListener() InstrumentedFeedSlot.subscribe__listener(listener) - with mock.patch("scrapy.extensions.feedexport._FeedSlot", InstrumentedFeedSlot): + with mock.patch("scrapy.extensions.feedexport.FeedSlot", InstrumentedFeedSlot): _ = yield self.exported_data(items, settings) self.assertFalse(listener.start_without_finish) self.assertFalse(listener.finish_without_start) @@ -1003,7 +1004,7 @@ class FeedExportTest(FeedExportTestBase): listener = IsExportingListener() InstrumentedFeedSlot.subscribe__listener(listener) - with mock.patch("scrapy.extensions.feedexport._FeedSlot", InstrumentedFeedSlot): + with mock.patch("scrapy.extensions.feedexport.FeedSlot", InstrumentedFeedSlot): _ = yield self.exported_data(items, settings) self.assertFalse(listener.start_without_finish) self.assertFalse(listener.finish_without_start) @@ -1022,7 +1023,7 @@ class FeedExportTest(FeedExportTestBase): listener = IsExportingListener() InstrumentedFeedSlot.subscribe__listener(listener) - with mock.patch("scrapy.extensions.feedexport._FeedSlot", InstrumentedFeedSlot): + with mock.patch("scrapy.extensions.feedexport.FeedSlot", InstrumentedFeedSlot): _ = yield self.exported_data(items, settings) self.assertFalse(listener.start_without_finish) self.assertFalse(listener.finish_without_start) @@ -2651,6 +2652,83 @@ class BatchDeliveriesTest(FeedExportTestBase): stub.assert_no_pending_responses() +# Test that the FeedExporer sends the feed_exporter_closed and feed_slot_closed signals +class FeedExporterSignalsTest(unittest.TestCase): + items = [ + {"foo": "bar1", "egg": "spam1"}, + {"foo": "bar2", "egg": "spam2", "baz": "quux2"}, + {"foo": "bar3", "baz": "quux3"}, + ] + + with tempfile.NamedTemporaryFile(suffix="json") as tmp: + settings = { + "FEEDS": { + f"file:///{tmp.name}": { + "format": "json", + }, + }, + } + + def feed_exporter_closed_signal_handler(self): + self.feed_exporter_closed_received = True + + 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 + + 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 + + def run_signaled_feed_exporter( + self, feed_exporter_signal_handler, feed_slot_signal_handler + ): + crawler = get_crawler(settings_dict=self.settings) + feed_exporter = FeedExporter.from_crawler(crawler) + spider = scrapy.Spider("default") + spider.crawler = crawler + crawler.signals.connect( + feed_exporter_signal_handler, + signal=signals.feed_exporter_closed, + ) + crawler.signals.connect( + feed_slot_signal_handler, signal=signals.feed_slot_closed + ) + feed_exporter.open_spider(spider) + for item in self.items: + feed_exporter.item_scraped(item, spider) + defer.ensureDeferred(feed_exporter.close_spider(spider)) + + def test_feed_exporter_signals_sent(self): + self.feed_exporter_closed_received = False + self.feed_slot_closed_received = False + + self.run_signaled_feed_exporter( + self.feed_exporter_closed_signal_handler, + self.feed_slot_closed_signal_handler, + ) + self.assertTrue(self.feed_slot_closed_received) + self.assertTrue(self.feed_exporter_closed_received) + + def test_feed_exporter_signals_sent_deferred(self): + 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, + ) + self.assertTrue(self.feed_slot_closed_received) + self.assertTrue(self.feed_exporter_closed_received) + + class FeedExportInitTest(unittest.TestCase): def test_unsupported_storage(self): settings = {