mirror of https://github.com/scrapy/scrapy.git
Merge branch 'master' into flake8-remove-e128
This commit is contained in:
commit
25e9bc2d0d
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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.",
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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:
|
||||
|
|
|
|||
|
|
@ -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):
|
||||
|
|
|
|||
|
|
@ -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()
|
||||
|
|
|
|||
|
|
@ -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(
|
||||
|
|
|
|||
|
|
@ -730,6 +730,9 @@ class Http11ProxyTestCase(HttpProxyTestCase):
|
|||
|
||||
class HttpDownloadHandlerMock:
|
||||
|
||||
def __init__(self, *args, **kwargs):
|
||||
pass
|
||||
|
||||
def download_request(self, request, spider):
|
||||
return request
|
||||
|
||||
|
|
|
|||
|
|
@ -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<html>\n"
|
||||
b" <head>\n"
|
||||
b" <meta http-equiv=\"refresh\" content=\"0;URL=/redirected\">\n"
|
||||
b" </head>\n"
|
||||
b" <body bgcolor=\"#FFFFFF\" text=\"#000000\">\n"
|
||||
b" <a href=\"/redirected\">click here</a>\n"
|
||||
b" </body>\n"
|
||||
b"</html>\n"
|
||||
)
|
||||
elif self.run.getpath(request.url) == "/tem999.html":
|
||||
self.assertEqual(
|
||||
joined_data,
|
||||
b"\n<html>\n"
|
||||
b" <head><title>404 - No Such Resource</title></head>\n"
|
||||
b" <body>\n"
|
||||
b" <h1>No Such Resource</h1>\n"
|
||||
b" <p>File not found.</p>\n"
|
||||
b" </body>\n"
|
||||
b"</html>\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):
|
||||
|
|
|
|||
|
|
@ -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='<html><body><p>some text</p></body></html>')
|
||||
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='<html><body><p>some text</p></body></html>')
|
||||
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)
|
||||
|
|
|
|||
Loading…
Reference in New Issue