diff --git a/docs/topics/signals.rst b/docs/topics/signals.rst index 8661f86a0..7fe63a7b0 100644 --- a/docs/topics/signals.rst +++ b/docs/topics/signals.rst @@ -112,7 +112,7 @@ engine_started Sent when the Scrapy engine has started crawling. - This signal supports returning deferreds from their handlers. + This signal supports returning deferreds from its 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 @@ -127,7 +127,7 @@ engine_stopped Sent when the Scrapy engine is stopped (for example, when a crawling process has finished). - This signal supports returning deferreds from their handlers. + This signal supports returning deferreds from its handlers. Item signals ------------ @@ -149,7 +149,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 their handlers. + This signal supports returning deferreds from its handlers. :param item: the item scraped :type item: dict or :class:`~scrapy.item.Item` object @@ -169,7 +169,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 their handlers. + This signal supports returning deferreds from its handlers. :param item: the item dropped from the :ref:`topics-item-pipeline` :type item: dict or :class:`~scrapy.item.Item` object @@ -194,7 +194,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 their handlers. + This signal supports returning deferreds from its handlers. :param item: the item dropped from the :ref:`topics-item-pipeline` :type item: dict or :class:`~scrapy.item.Item` object @@ -220,7 +220,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 their handlers. + This signal supports returning deferreds from its handlers. :param spider: the spider which has been closed :type spider: :class:`~scrapy.spiders.Spider` object @@ -244,7 +244,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 their handlers. + This signal supports returning deferreds from its handlers. :param spider: the spider which has been opened :type spider: :class:`~scrapy.spiders.Spider` object @@ -268,7 +268,7 @@ spider_idle You may raise a :exc:`~scrapy.exceptions.DontCloseSpider` exception to prevent the spider from being closed. - This signal does not support returning deferreds from their handlers. + This signal does not support returning deferreds from its handlers. :param spider: the spider which has gone idle :type spider: :class:`~scrapy.spiders.Spider` object @@ -287,7 +287,7 @@ spider_error Sent when a spider callback generates an error (i.e. raises an exception). - This signal does not support returning deferreds from their handlers. + This signal does not support returning deferreds from its handlers. :param failure: the exception raised :type failure: twisted.python.failure.Failure @@ -310,7 +310,7 @@ request_scheduled Sent when the engine schedules a :class:`~scrapy.http.Request`, to be downloaded later. - The signal does not support returning deferreds from their handlers. + This signal does not support returning deferreds from its handlers. :param request: the request that reached the scheduler :type request: :class:`~scrapy.http.Request` object @@ -327,7 +327,7 @@ request_dropped Sent when a :class:`~scrapy.http.Request`, scheduled by the engine to be downloaded later, is rejected by the scheduler. - The signal does not support returning deferreds from their handlers. + This signal does not support returning deferreds from its handlers. :param request: the request that reached the scheduler :type request: :class:`~scrapy.http.Request` object @@ -343,7 +343,7 @@ request_reached_downloader Sent when a :class:`~scrapy.http.Request` reached downloader. - The signal does not support returning deferreds from their handlers. + This signal does not support returning deferreds from its handlers. :param request: the request that reached downloader :type request: :class:`~scrapy.http.Request` object @@ -370,6 +370,29 @@ request_left_downloader :param spider: the spider that yielded the request :type spider: :class:`~scrapy.spiders.Spider` object +bytes_received +~~~~~~~~~~~~~~ + +.. signal:: bytes_received +.. function:: bytes_received(data, request, spider) + + Sent by the HTTP 1.1 and S3 download handlers when a group of bytes is + received for a specific request. This signal might be fired multiple + times for the same request, with partial data each time. For instance, + a possible scenario for a 25 kb response would be two signals fired + with 10 kb of data, and a final one with 5 kb of data. + + This signal does not support returning deferreds from its handlers. + + :param data: the data received by the download handler + :type spider: :class:`bytes` object + + :param request: the request that generated the response + :type request: :class:`~scrapy.http.Request` object + + :param spider: the spider associated with the response + :type spider: :class:`~scrapy.spiders.Spider` object + Response signals ---------------- @@ -382,7 +405,7 @@ response_received Sent when the engine receives a new :class:`~scrapy.http.Response` from the downloader. - This signal does not support returning deferreds from their handlers. + This signal does not support returning deferreds from its handlers. :param response: the response received :type response: :class:`~scrapy.http.Response` object @@ -401,7 +424,7 @@ response_downloaded Sent by the downloader right after a ``HTTPResponse`` is downloaded. - This signal does not support returning deferreds from their handlers. + This signal does not support returning deferreds from its handlers. :param response: the response downloaded :type response: :class:`~scrapy.http.Response` object diff --git a/pytest.ini b/pytest.ini index f9769467d..6e0287011 100644 --- a/pytest.ini +++ b/pytest.ini @@ -164,6 +164,7 @@ flake8-ignore = scrapy/shell.py E501 scrapy/signalmanager.py E501 scrapy/spiderloader.py F841 E501 + scrapy/squeues.py E501 scrapy/statscollectors.py E501 # tests tests/__init__.py E402 E501 diff --git a/scrapy/core/downloader/handlers/http11.py b/scrapy/core/downloader/handlers/http11.py index 9be8ffdfb..c21491f52 100644 --- a/scrapy/core/downloader/handlers/http11.py +++ b/scrapy/core/downloader/handlers/http11.py @@ -18,6 +18,7 @@ from twisted.web.http_headers import Headers as TxHeaders from twisted.web.iweb import IBodyProducer, UNKNOWN_LENGTH from zope.interface import implementer +from scrapy import signals from scrapy.core.downloader.tls import openssl_methods from scrapy.core.downloader.webclient import _parse from scrapy.exceptions import ScrapyDeprecationWarning @@ -34,6 +35,8 @@ class HTTP11DownloadHandler: lazy = False def __init__(self, settings, crawler=None): + self._crawler = crawler + from twisted.internet import reactor self._pool = HTTPConnectionPool(reactor, persistent=True) self._pool.maxPersistentPerHost = settings.getint('CONCURRENT_REQUESTS_PER_DOMAIN') @@ -79,6 +82,7 @@ class HTTP11DownloadHandler: maxsize=getattr(spider, 'download_maxsize', self._default_maxsize), warnsize=getattr(spider, 'download_warnsize', self._default_warnsize), fail_on_dataloss=self._fail_on_dataloss, + crawler=self._crawler, ) return agent.download_request(request) @@ -276,7 +280,7 @@ class ScrapyAgent: _TunnelingAgent = TunnelingAgent def __init__(self, contextFactory=None, connectTimeout=10, bindAddress=None, pool=None, - maxsize=0, warnsize=0, fail_on_dataloss=True): + maxsize=0, warnsize=0, fail_on_dataloss=True, crawler=None): self._contextFactory = contextFactory self._connectTimeout = connectTimeout self._bindAddress = bindAddress @@ -285,6 +289,7 @@ class ScrapyAgent: self._warnsize = warnsize self._fail_on_dataloss = fail_on_dataloss self._txresponse = None + self._crawler = crawler def _get_agent(self, request, timeout): from twisted.internet import reactor @@ -407,7 +412,15 @@ class ScrapyAgent: d = defer.Deferred(_cancel) txresponse.deliverBody( - _ResponseReader(d, txresponse, request, maxsize, warnsize, fail_on_dataloss) + _ResponseReader( + finished=d, + txresponse=txresponse, + request=request, + maxsize=maxsize, + warnsize=warnsize, + fail_on_dataloss=fail_on_dataloss, + crawler=self._crawler, + ) ) # save response for timeouts @@ -449,7 +462,7 @@ class _RequestBodyProducer: class _ResponseReader(protocol.Protocol): - def __init__(self, finished, txresponse, request, maxsize, warnsize, fail_on_dataloss): + def __init__(self, finished, txresponse, request, maxsize, warnsize, fail_on_dataloss, crawler): self._finished = finished self._txresponse = txresponse self._request = request @@ -462,6 +475,7 @@ class _ResponseReader(protocol.Protocol): self._bytes_received = 0 self._certificate = None self._ip_address = None + self._crawler = crawler def connectionMade(self): if self._certificate is None: @@ -479,6 +493,13 @@ class _ResponseReader(protocol.Protocol): self._bodybuf.write(bodyBytes) self._bytes_received += len(bodyBytes) + self._crawler.signals.send_catch_log( + signal=signals.bytes_received, + data=bodyBytes, + request=self._request, + spider=self._crawler.spider, + ) + if self._maxsize and self._bytes_received > self._maxsize: logger.error("Received (%(bytes)s) bytes larger than download " "max size (%(maxsize)s) in request %(request)s.", diff --git a/scrapy/exporters.py b/scrapy/exporters.py index 0cb6cef98..349a9586b 100644 --- a/scrapy/exporters.py +++ b/scrapy/exporters.py @@ -250,7 +250,7 @@ class CsvItemExporter(BaseItemExporter): class PickleItemExporter(BaseItemExporter): - def __init__(self, file, protocol=2, **kwargs): + def __init__(self, file, protocol=4, **kwargs): super().__init__(**kwargs) self.file = file self.protocol = protocol diff --git a/scrapy/extensions/httpcache.py b/scrapy/extensions/httpcache.py index 6289efec0..6294a9b52 100644 --- a/scrapy/extensions/httpcache.py +++ b/scrapy/extensions/httpcache.py @@ -251,7 +251,7 @@ class DbmCacheStorage: 'headers': dict(response.headers), 'body': response.body, } - self.db['%s_data' % key] = pickle.dumps(data, protocol=2) + self.db['%s_data' % key] = pickle.dumps(data, protocol=4) self.db['%s_time' % key] = str(time()) def _read_data(self, spider, request): @@ -318,7 +318,7 @@ class FilesystemCacheStorage: with self._open(os.path.join(rpath, 'meta'), 'wb') as f: f.write(to_bytes(repr(metadata))) with self._open(os.path.join(rpath, 'pickled_meta'), 'wb') as f: - pickle.dump(metadata, f, protocol=2) + pickle.dump(metadata, f, protocol=4) with self._open(os.path.join(rpath, 'response_headers'), 'wb') as f: f.write(headers_dict_to_raw(response.headers)) with self._open(os.path.join(rpath, 'response_body'), 'wb') as f: diff --git a/scrapy/extensions/spiderstate.py b/scrapy/extensions/spiderstate.py index 2e5ff569f..bea00596e 100644 --- a/scrapy/extensions/spiderstate.py +++ b/scrapy/extensions/spiderstate.py @@ -26,7 +26,7 @@ class SpiderState: def spider_closed(self, spider): if self.jobdir: with open(self.statefn, 'wb') as f: - pickle.dump(spider.state, f, protocol=2) + pickle.dump(spider.state, f, protocol=4) def spider_opened(self, spider): if self.jobdir and os.path.exists(self.statefn): diff --git a/scrapy/signals.py b/scrapy/signals.py index cd7ed7fb1..c61ae6ec3 100644 --- a/scrapy/signals.py +++ b/scrapy/signals.py @@ -17,6 +17,7 @@ request_reached_downloader = object() request_left_downloader = object() response_received = object() response_downloaded = object() +bytes_received = object() item_scraped = object() item_dropped = object() item_error = object() diff --git a/scrapy/squeues.py b/scrapy/squeues.py index d0686dac3..c7ad4d53d 100644 --- a/scrapy/squeues.py +++ b/scrapy/squeues.py @@ -81,12 +81,11 @@ def _scrapy_non_serialization_queue(queue_class): def _pickle_serialize(obj): try: - return pickle.dumps(obj, protocol=2) - # Python <= 3.4 raises pickle.PicklingError here while - # 3.5 <= Python < 3.6 raises AttributeError and - # Python >= 3.6 raises TypeError + return pickle.dumps(obj, protocol=4) + # Both pickle.PicklingError and AttributeError can be raised by pickle.dump(s) + # TypeError is raised from parsel.Selector except (pickle.PicklingError, AttributeError, TypeError) as e: - raise ValueError(str(e)) + raise ValueError(str(e)) from e PickleFifoDiskQueueNonRequest = _serializable_queue( diff --git a/tests/test_downloader_handlers.py b/tests/test_downloader_handlers.py index 1a05b679a..c1e6f744b 100644 --- a/tests/test_downloader_handlers.py +++ b/tests/test_downloader_handlers.py @@ -730,6 +730,9 @@ class Http11ProxyTestCase(HttpProxyTestCase): class HttpDownloadHandlerMock: + def __init__(self, *args, **kwargs): + pass + def download_request(self, request, spider): return request diff --git a/tests/test_engine.py b/tests/test_engine.py index 5b7a4e676..acfe94f63 100644 --- a/tests/test_engine.py +++ b/tests/test_engine.py @@ -13,22 +13,24 @@ module with the ``runserver`` argument:: import os import re import sys +from collections import defaultdict from urllib.parse import urlparse from twisted.internet import reactor, defer -from twisted.web import server, static, util from twisted.trial import unittest +from twisted.web import server, static, util +from pydispatch import dispatcher from scrapy import signals from scrapy.core.engine import ExecutionEngine -from scrapy.utils.test import get_crawler -from pydispatch import dispatcher -from tests import tests_datadir -from scrapy.spiders import Spider +from scrapy.http import Request from scrapy.item import Item, Field from scrapy.linkextractors import LinkExtractor -from scrapy.http import Request +from scrapy.spiders import Spider from scrapy.utils.signal import disconnect_all +from scrapy.utils.test import get_crawler + +from tests import tests_datadir, get_testdata class TestItem(Item): @@ -88,6 +90,8 @@ def start_test_site(debug=False): r = static.File(root_dir) r.putChild(b"redirect", util.Redirect(b"/redirected")) r.putChild(b"redirected", static.Data(b"Redirected here", "text/plain")) + numbers = [str(x).encode("utf8") for x in range(2**14)] + r.putChild(b"numbers", static.Data(b"".join(numbers), "text/plain")) port = reactor.listenTCP(0, server.Site(r), interface="127.0.0.1") if debug: @@ -107,15 +111,20 @@ class CrawlerRun: self.reqreached = [] self.itemerror = [] self.itemresp = [] - self.signals_catched = {} + self.bytes = defaultdict(lambda: list()) + self.signals_caught = {} self.spider_class = spider_class def run(self): self.port = start_test_site() self.portno = self.port.getHost().port - start_urls = [self.geturl("/"), self.geturl("/redirect"), - self.geturl("/redirect")] # a duplicate + start_urls = [ + self.geturl("/"), + self.geturl("/redirect"), + self.geturl("/redirect"), # duplicate + self.geturl("/numbers"), + ] for name, signal in vars(signals).items(): if not name.startswith('_'): @@ -124,6 +133,7 @@ class CrawlerRun: self.crawler = get_crawler(self.spider_class) self.crawler.signals.connect(self.item_scraped, signals.item_scraped) self.crawler.signals.connect(self.item_error, signals.item_error) + self.crawler.signals.connect(self.bytes_received, signals.bytes_received) self.crawler.signals.connect(self.request_scheduled, signals.request_scheduled) self.crawler.signals.connect(self.request_dropped, signals.request_dropped) self.crawler.signals.connect(self.request_reached, signals.request_reached_downloader) @@ -155,6 +165,9 @@ class CrawlerRun: def item_scraped(self, item, spider, response): self.itemresp.append((item, response)) + def bytes_received(self, data, request, spider): + self.bytes[request].append(data) + def request_scheduled(self, request, spider): self.reqplug.append((request, spider)) @@ -172,7 +185,7 @@ class CrawlerRun: signalargs = kwargs.copy() sig = signalargs.pop('signal') signalargs.pop('sender', None) - self.signals_catched[sig] = signalargs + self.signals_caught[sig] = signalargs class EngineTest(unittest.TestCase): @@ -183,16 +196,17 @@ class EngineTest(unittest.TestCase): self.run = CrawlerRun(spider) yield self.run.run() self._assert_visited_urls() - self._assert_scheduled_requests(urls_to_visit=8) + self._assert_scheduled_requests(urls_to_visit=9) self._assert_downloaded_responses() self._assert_scraped_items() - self._assert_signals_catched() + self._assert_signals_caught() + self._assert_bytes_received() @defer.inlineCallbacks def test_crawler_dupefilter(self): self.run = CrawlerRun(TestDupeFilterSpider) yield self.run.run() - self._assert_scheduled_requests(urls_to_visit=7) + self._assert_scheduled_requests(urls_to_visit=8) self._assert_dropped_requests() @defer.inlineCallbacks @@ -229,8 +243,8 @@ class EngineTest(unittest.TestCase): def _assert_downloaded_responses(self): # response tests - self.assertEqual(8, len(self.run.respplug)) - self.assertEqual(8, len(self.run.reqreached)) + self.assertEqual(9, len(self.run.respplug)) + self.assertEqual(9, len(self.run.reqreached)) for response, _ in self.run.respplug: if self.run.getpath(response.url) == '/item999.html': @@ -263,19 +277,61 @@ class EngineTest(unittest.TestCase): self.assertEqual('Item 2 name', item['name']) self.assertEqual('200', item['price']) - def _assert_signals_catched(self): - assert signals.engine_started in self.run.signals_catched - assert signals.engine_stopped in self.run.signals_catched - assert signals.spider_opened in self.run.signals_catched - assert signals.spider_idle in self.run.signals_catched - assert signals.spider_closed in self.run.signals_catched + def _assert_bytes_received(self): + self.assertEqual(9, len(self.run.bytes)) + for request, data in self.run.bytes.items(): + joined_data = b"".join(data) + if self.run.getpath(request.url) == "/": + self.assertEqual(joined_data, get_testdata("test_site", "index.html")) + elif self.run.getpath(request.url) == "/item1.html": + self.assertEqual(joined_data, get_testdata("test_site", "item1.html")) + elif self.run.getpath(request.url) == "/item2.html": + self.assertEqual(joined_data, get_testdata("test_site", "item2.html")) + elif self.run.getpath(request.url) == "/redirected": + self.assertEqual(joined_data, b"Redirected here") + elif self.run.getpath(request.url) == '/redirect': + self.assertEqual( + joined_data, + b"\n\n" + b" \n" + b" \n" + b" \n" + b" \n" + b" click here\n" + b" \n" + b"\n" + ) + elif self.run.getpath(request.url) == "/tem999.html": + self.assertEqual( + joined_data, + b"\n\n" + b" 404 - No Such Resource\n" + b" \n" + b"

No Such Resource

\n" + b"

File not found.

\n" + b" \n" + b"\n" + ) + elif self.run.getpath(request.url) == "/numbers": + # signal was fired multiple times + self.assertTrue(len(data) > 1) + # bytes were received in order + numbers = [str(x).encode("utf8") for x in range(2**14)] + self.assertEqual(joined_data, b"".join(numbers)) + + def _assert_signals_caught(self): + assert signals.engine_started in self.run.signals_caught + assert signals.engine_stopped in self.run.signals_caught + assert signals.spider_opened in self.run.signals_caught + assert signals.spider_idle in self.run.signals_caught + assert signals.spider_closed in self.run.signals_caught self.assertEqual({'spider': self.run.spider}, - self.run.signals_catched[signals.spider_opened]) + self.run.signals_caught[signals.spider_opened]) self.assertEqual({'spider': self.run.spider}, - self.run.signals_catched[signals.spider_idle]) + self.run.signals_caught[signals.spider_idle]) self.assertEqual({'spider': self.run.spider, 'reason': 'finished'}, - self.run.signals_catched[signals.spider_closed]) + self.run.signals_caught[signals.spider_closed]) @defer.inlineCallbacks def test_close_downloader(self): diff --git a/tests/test_squeues.py b/tests/test_squeues.py index 5ad8035f7..d2cf9135f 100644 --- a/tests/test_squeues.py +++ b/tests/test_squeues.py @@ -1,4 +1,5 @@ import pickle +import sys from queuelib.tests import test_queue as t from scrapy.squeues import ( @@ -28,31 +29,13 @@ class TestLoader(ItemLoader): def nonserializable_object_test(self): q = self.queue() - try: - pickle.dumps(lambda x: x) - except Exception: - # Trigger Twisted bug #7989 - import twisted.persisted.styles # NOQA - self.assertRaises(ValueError, q.push, lambda x: x) - else: - # Use a different unpickleable object - class A: - pass - - a = A() - a.__reduce__ = a.__reduce_ex__ = None - self.assertRaises(ValueError, q.push, a) + self.assertRaises(ValueError, q.push, lambda x: x) # Selectors should fail (lxml.html.HtmlElement objects can't be pickled) sel = Selector(text='

some text

') self.assertRaises(ValueError, q.push, sel) -class MarshalFifoDiskQueueTest(t.FifoDiskQueueTest): - - chunksize = 100000 - - def queue(self): - return MarshalFifoDiskQueue(self.qpath, chunksize=self.chunksize) +class FifoDiskQueueTestMixin: def test_serialize(self): q = self.queue() @@ -66,6 +49,13 @@ class MarshalFifoDiskQueueTest(t.FifoDiskQueueTest): test_nonserializable_object = nonserializable_object_test +class MarshalFifoDiskQueueTest(t.FifoDiskQueueTest, FifoDiskQueueTestMixin): + chunksize = 100000 + + def queue(self): + return MarshalFifoDiskQueue(self.qpath, chunksize=self.chunksize) + + class ChunkSize1MarshalFifoDiskQueueTest(MarshalFifoDiskQueueTest): chunksize = 1 @@ -82,7 +72,7 @@ class ChunkSize4MarshalFifoDiskQueueTest(MarshalFifoDiskQueueTest): chunksize = 4 -class PickleFifoDiskQueueTest(MarshalFifoDiskQueueTest): +class PickleFifoDiskQueueTest(t.FifoDiskQueueTest, FifoDiskQueueTestMixin): chunksize = 100000 @@ -116,6 +106,21 @@ class PickleFifoDiskQueueTest(MarshalFifoDiskQueueTest): self.assertEqual(r.url, r2.url) assert r2.meta['request'] is r2 + def test_non_pickable_object(self): + q = self.queue() + try: + q.push(lambda x: x) + except ValueError as exc: + if hasattr(sys, "pypy_version_info"): + self.assertIsInstance(exc.__context__, pickle.PicklingError) + else: + self.assertIsInstance(exc.__context__, AttributeError) + sel = Selector(text='

some text

') + try: + q.push(sel) + except ValueError as exc: + self.assertIsInstance(exc.__context__, TypeError) + class ChunkSize1PickleFifoDiskQueueTest(PickleFifoDiskQueueTest): chunksize = 1 @@ -133,10 +138,7 @@ class ChunkSize4PickleFifoDiskQueueTest(PickleFifoDiskQueueTest): chunksize = 4 -class MarshalLifoDiskQueueTest(t.LifoDiskQueueTest): - - def queue(self): - return MarshalLifoDiskQueue(self.qpath) +class LifoDiskQueueTestMixin: def test_serialize(self): q = self.queue() @@ -150,7 +152,13 @@ class MarshalLifoDiskQueueTest(t.LifoDiskQueueTest): test_nonserializable_object = nonserializable_object_test -class PickleLifoDiskQueueTest(MarshalLifoDiskQueueTest): +class MarshalLifoDiskQueueTest(t.LifoDiskQueueTest, LifoDiskQueueTestMixin): + + def queue(self): + return MarshalLifoDiskQueue(self.qpath) + + +class PickleLifoDiskQueueTest(t.LifoDiskQueueTest, LifoDiskQueueTestMixin): def queue(self): return PickleLifoDiskQueue(self.qpath)