From b247fa9982a390e1380c46e390069c6058d97921 Mon Sep 17 00:00:00 2001 From: Ricardo Amendoeira Date: Mon, 29 Mar 2021 01:48:28 +0100 Subject: [PATCH 01/16] Include loading settings in `Running multiple spiders in the same process` section The example in the documentation doesn't take into account the project settings --- docs/topics/practices.rst | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/docs/topics/practices.rst b/docs/topics/practices.rst index cf1de1bd1..db1ed362e 100644 --- a/docs/topics/practices.rst +++ b/docs/topics/practices.rst @@ -118,6 +118,7 @@ Here is an example that runs multiple spiders simultaneously: :: import scrapy + from scrapy.utils.project import get_project_settings from scrapy.crawler import CrawlerProcess class MySpider1(scrapy.Spider): @@ -128,7 +129,8 @@ Here is an example that runs multiple spiders simultaneously: # Your second spider definition ... - process = CrawlerProcess() + settings = get_project_settings() + process = CrawlerProcess(settings) process.crawl(MySpider1) process.crawl(MySpider2) process.start() # the script will block here until all crawling jobs are finished From e7d51886ef90f5b2d4fa13911381680ac192fc37 Mon Sep 17 00:00:00 2001 From: Mayank Singhal <17mayanksinghal@gmail.com> Date: Tue, 6 Apr 2021 02:21:18 +0530 Subject: [PATCH 02/16] Find bash from PATH instead of /bin/bash --- docs/Makefile | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/docs/Makefile b/docs/Makefile index ff68bf1ae..87d5d3047 100644 --- a/docs/Makefile +++ b/docs/Makefile @@ -8,7 +8,7 @@ PYTHON = python SPHINXOPTS = PAPER = SOURCES = -SHELL = /bin/bash +SHELL = /usr/bin/env bash ALLSPHINXOPTS = -b $(BUILDER) -d build/doctrees \ -D latex_elements.papersize=$(PAPER) \ From 8603f9d7a5524b7709b4e8a8fc04a75a9a4f0ffe Mon Sep 17 00:00:00 2001 From: Ricardo Amendoeira Date: Tue, 6 Apr 2021 20:23:07 +0100 Subject: [PATCH 03/16] Apply changes to other examples in the same section. --- docs/topics/practices.rst | 10 +++++++--- 1 file changed, 7 insertions(+), 3 deletions(-) diff --git a/docs/topics/practices.rst b/docs/topics/practices.rst index db1ed362e..15ac520e2 100644 --- a/docs/topics/practices.rst +++ b/docs/topics/practices.rst @@ -118,8 +118,8 @@ Here is an example that runs multiple spiders simultaneously: :: import scrapy - from scrapy.utils.project import get_project_settings from scrapy.crawler import CrawlerProcess + from scrapy.utils.project import get_project_settings class MySpider1(scrapy.Spider): # Your first spider definition @@ -143,6 +143,7 @@ Same example using :class:`~scrapy.crawler.CrawlerRunner`: from twisted.internet import reactor from scrapy.crawler import CrawlerRunner from scrapy.utils.log import configure_logging + from scrapy.utils.project import get_project_settings class MySpider1(scrapy.Spider): # Your first spider definition @@ -153,7 +154,8 @@ Same example using :class:`~scrapy.crawler.CrawlerRunner`: ... configure_logging() - runner = CrawlerRunner() + settings = get_project_settings() + runner = CrawlerRunner(settings) runner.crawl(MySpider1) runner.crawl(MySpider2) d = runner.join() @@ -168,6 +170,7 @@ Same example but running the spiders sequentially by chaining the deferreds: from twisted.internet import reactor, defer from scrapy.crawler import CrawlerRunner from scrapy.utils.log import configure_logging + from scrapy.utils.project import get_project_settings class MySpider1(scrapy.Spider): # Your first spider definition @@ -178,7 +181,8 @@ Same example but running the spiders sequentially by chaining the deferreds: ... configure_logging() - runner = CrawlerRunner() + settings = get_project_settings() + runner = CrawlerRunner(settings) @defer.inlineCallbacks def crawl(): From 7e23677b52b659b11471a63f3be9905a0bbaf995 Mon Sep 17 00:00:00 2001 From: Eugenio Lacuesta <1731933+elacuesta@users.noreply.github.com> Date: Tue, 20 Apr 2021 08:45:28 -0300 Subject: [PATCH 04/16] Engine: deprecations and type hints (#5090) --- docs/topics/telnetconsole.rst | 3 +- scrapy/core/engine.py | 419 ++++++++++--------- scrapy/core/scraper.py | 2 +- scrapy/downloadermiddlewares/robotstxt.py | 2 +- scrapy/extensions/memusage.py | 6 +- scrapy/pipelines/media.py | 2 +- scrapy/shell.py | 2 +- scrapy/utils/engine.py | 3 +- tests/test_downloadermiddleware_robotstxt.py | 12 +- tests/test_engine.py | 110 ++++- 10 files changed, 336 insertions(+), 225 deletions(-) diff --git a/docs/topics/telnetconsole.rst b/docs/topics/telnetconsole.rst index 9802a34a2..832829b75 100644 --- a/docs/topics/telnetconsole.rst +++ b/docs/topics/telnetconsole.rst @@ -110,11 +110,10 @@ using the telnet console:: Execution engine status time()-engine.start_time : 8.62972998619 - engine.has_capacity() : False len(engine.downloader.active) : 16 engine.scraper.is_idle() : False engine.spider.name : followall - engine.spider_is_idle(engine.spider) : False + engine.spider_is_idle() : False engine.slot.closing : False len(engine.slot.inprogress) : 16 len(engine.slot.scheduler.dqs or []) : 0 diff --git a/scrapy/core/engine.py b/scrapy/core/engine.py index 93bcdb49a..edfac87c6 100644 --- a/scrapy/core/engine.py +++ b/scrapy/core/engine.py @@ -1,51 +1,61 @@ """ -This is the Scrapy engine which controls the Scheduler, Downloader and Spiders. +This is the Scrapy engine which controls the Scheduler, Downloader and Spider. For more information see docs/topics/architecture.rst """ import logging +import warnings from time import time +from typing import Callable, Iterable, Iterator, Optional, Set, Union -from twisted.internet import defer, task +from twisted.internet.defer import Deferred, inlineCallbacks, succeed +from twisted.internet.task import LoopingCall from twisted.python.failure import Failure from scrapy import signals from scrapy.core.scraper import Scraper -from scrapy.exceptions import DontCloseSpider +from scrapy.exceptions import DontCloseSpider, ScrapyDeprecationWarning from scrapy.http import Response, Request -from scrapy.utils.misc import load_object -from scrapy.utils.reactor import CallLaterOnce +from scrapy.spiders import Spider from scrapy.utils.log import logformatter_adapter, failure_to_exc_info +from scrapy.utils.misc import create_instance, load_object +from scrapy.utils.reactor import CallLaterOnce + logger = logging.getLogger(__name__) class Slot: - - def __init__(self, start_requests, close_if_idle, nextcall, scheduler): - self.closing = False - self.inprogress = set() # requests in progress - self.start_requests = iter(start_requests) + def __init__( + self, + start_requests: Iterable, + close_if_idle: bool, + nextcall: CallLaterOnce, + scheduler, + ) -> None: + self.closing: Optional[Deferred] = None + self.inprogress: Set[Request] = set() + self.start_requests: Optional[Iterator] = iter(start_requests) self.close_if_idle = close_if_idle self.nextcall = nextcall self.scheduler = scheduler - self.heartbeat = task.LoopingCall(nextcall.schedule) + self.heartbeat = LoopingCall(nextcall.schedule) - def add_request(self, request): + def add_request(self, request: Request) -> None: self.inprogress.add(request) - def remove_request(self, request): + def remove_request(self, request: Request) -> None: self.inprogress.remove(request) self._maybe_fire_closing() - def close(self): - self.closing = defer.Deferred() + def close(self) -> Deferred: + self.closing = Deferred() self._maybe_fire_closing() return self.closing - def _maybe_fire_closing(self): - if self.closing and not self.inprogress: + def _maybe_fire_closing(self) -> None: + if self.closing is not None and not self.inprogress: if self.nextcall: self.nextcall.cancel() if self.heartbeat.running: @@ -54,210 +64,224 @@ class Slot: class ExecutionEngine: - - def __init__(self, crawler, spider_closed_callback): + def __init__(self, crawler, spider_closed_callback: Callable) -> None: self.crawler = crawler self.settings = crawler.settings self.signals = crawler.signals self.logformatter = crawler.logformatter - self.slot = None - self.spider = None + self.slot: Optional[Slot] = None + self.spider: Optional[Spider] = None self.running = False self.paused = False - self.scheduler_cls = load_object(self.settings['SCHEDULER']) + self.scheduler_cls = load_object(crawler.settings["SCHEDULER"]) downloader_cls = load_object(self.settings['DOWNLOADER']) self.downloader = downloader_cls(crawler) self.scraper = Scraper(crawler) self._spider_closed_callback = spider_closed_callback - @defer.inlineCallbacks - def start(self): - """Start the execution engine""" + @inlineCallbacks + def start(self) -> Deferred: if self.running: raise RuntimeError("Engine already running") self.start_time = time() yield self.signals.send_catch_log_deferred(signal=signals.engine_started) self.running = True - self._closewait = defer.Deferred() + self._closewait = Deferred() yield self._closewait - def stop(self): - """Stop the execution engine gracefully""" + def stop(self) -> Deferred: + """Gracefully stop the execution engine""" + @inlineCallbacks + def _finish_stopping_engine(_) -> Deferred: + yield self.signals.send_catch_log_deferred(signal=signals.engine_stopped) + self._closewait.callback(None) + if not self.running: raise RuntimeError("Engine not running") + self.running = False - dfd = self._close_all_spiders() - return dfd.addBoth(lambda _: self._finish_stopping_engine()) + dfd = self.close_spider(self.spider, reason="shutdown") if self.spider is not None else succeed(None) + return dfd.addBoth(_finish_stopping_engine) - def close(self): - """Close the execution engine gracefully. - - If it has already been started, stop it. In all cases, close all spiders - and the downloader. + def close(self) -> Deferred: + """ + Gracefully close the execution engine. + If it has already been started, stop it. In all cases, close the spider and the downloader. """ if self.running: - # Will also close spiders and downloader - return self.stop() - elif self.open_spiders: - # Will also close downloader - return self._close_all_spiders() - else: - return defer.succeed(self.downloader.close()) + return self.stop() # will also close spider and downloader + if self.spider is not None: + return self.close_spider(self.spider, reason="shutdown") # will also close downloader + return succeed(self.downloader.close()) - def pause(self): - """Pause the execution engine""" + def pause(self) -> None: self.paused = True - def unpause(self): - """Resume the execution engine""" + def unpause(self) -> None: self.paused = False - def _next_request(self, spider): - slot = self.slot - if not slot: - return + def _next_request(self) -> None: + assert self.slot is not None # typing + assert self.spider is not None # typing if self.paused: - return + return None - while not self._needs_backout(spider): - if not self._next_request_from_scheduler(spider): - break + while not self._needs_backout() and self._next_request_from_scheduler() is not None: + pass - if slot.start_requests and not self._needs_backout(spider): + if self.slot.start_requests is not None and not self._needs_backout(): try: - request = next(slot.start_requests) + request = next(self.slot.start_requests) except StopIteration: - slot.start_requests = None + self.slot.start_requests = None except Exception: - slot.start_requests = None - logger.error('Error while obtaining start requests', - exc_info=True, extra={'spider': spider}) + self.slot.start_requests = None + logger.error('Error while obtaining start requests', exc_info=True, extra={'spider': self.spider}) else: - self.crawl(request, spider) + self.crawl(request) - if self.spider_is_idle(spider) and slot.close_if_idle: - self._spider_idle(spider) + if self.spider_is_idle() and self.slot.close_if_idle: + self._spider_idle() - def _needs_backout(self, spider): - slot = self.slot + def _needs_backout(self) -> bool: return ( not self.running - or slot.closing + or self.slot.closing # type: ignore[union-attr] or self.downloader.needs_backout() - or self.scraper.slot.needs_backout() + or self.scraper.slot.needs_backout() # type: ignore[union-attr] ) - def _next_request_from_scheduler(self, spider): - slot = self.slot - request = slot.scheduler.next_request() - if not request: - return - d = self._download(request, spider) - d.addBoth(self._handle_downloader_output, request, spider) + def _next_request_from_scheduler(self) -> Optional[Deferred]: + assert self.slot is not None # typing + assert self.spider is not None # typing + + request = self.slot.scheduler.next_request() + if request is None: + return None + + d = self._download(request, self.spider) + d.addBoth(self._handle_downloader_output, request, self.spider) d.addErrback(lambda f: logger.info('Error while handling downloader output', exc_info=failure_to_exc_info(f), - extra={'spider': spider})) - d.addBoth(lambda _: slot.remove_request(request)) + extra={'spider': self.spider})) + d.addBoth(lambda _: self.slot.remove_request(request)) d.addErrback(lambda f: logger.info('Error while removing request from slot', exc_info=failure_to_exc_info(f), - extra={'spider': spider})) - d.addBoth(lambda _: slot.nextcall.schedule()) + extra={'spider': self.spider})) + d.addBoth(lambda _: self.slot.nextcall.schedule()) d.addErrback(lambda f: logger.info('Error while scheduling new request', exc_info=failure_to_exc_info(f), - extra={'spider': spider})) + extra={'spider': self.spider})) return d - def _handle_downloader_output(self, response, request, spider): - if not isinstance(response, (Request, Response, Failure)): - raise TypeError( - "Incorrect type: expected Request, Response or Failure, got " - f"{type(response)}: {response!r}" - ) + def _handle_downloader_output( + self, result: Union[Request, Response, Failure], request: Request, spider: Spider + ) -> Optional[Deferred]: + if not isinstance(result, (Request, Response, Failure)): + raise TypeError(f"Incorrect type: expected Request, Response or Failure, got {type(result)}: {result!r}") + # downloader middleware can return requests (for example, redirects) - if isinstance(response, Request): - self.crawl(response, spider) - return - # response is a Response or Failure - d = self.scraper.enqueue_scrape(response, request, spider) - d.addErrback(lambda f: logger.error('Error while enqueuing downloader output', - exc_info=failure_to_exc_info(f), - extra={'spider': spider})) + if isinstance(result, Request): + self.crawl(result) + return None + + d = self.scraper.enqueue_scrape(result, request, spider) + d.addErrback( + lambda f: logger.error( + "Error while enqueuing downloader output", + exc_info=failure_to_exc_info(f), + extra={'spider': spider}, + ) + ) return d - def spider_is_idle(self, spider): - if not self.scraper.slot.is_idle(): - # scraper is not idle + def spider_is_idle(self, spider: Optional[Spider] = None) -> bool: + if spider is not None: + warnings.warn( + "Passing a 'spider' argument to ExecutionEngine.spider_is_idle is deprecated", + category=ScrapyDeprecationWarning, + stacklevel=2, + ) + if self.slot is None: + raise RuntimeError("Engine slot not assigned") + if not self.scraper.slot.is_idle(): # type: ignore[union-attr] return False - - if self.downloader.active: - # downloader has pending requests + if self.downloader.active: # downloader has pending requests return False - - if self.slot.start_requests is not None: - # not all start requests are handled + if self.slot.start_requests is not None: # not all start requests are handled return False - if self.slot.scheduler.has_pending_requests(): - # scheduler has pending requests return False - return True - @property - def open_spiders(self): - return [self.spider] if self.spider else [] + def crawl(self, request: Request, spider: Optional[Spider] = None) -> None: + """Inject the request into the spider <-> downloader pipeline""" + if spider is not None: + warnings.warn( + "Passing a 'spider' argument to ExecutionEngine.crawl is deprecated", + category=ScrapyDeprecationWarning, + stacklevel=2, + ) + if spider is not self.spider: + raise RuntimeError(f"The spider {spider.name!r} does not match the open spider") + if self.spider is None: + raise RuntimeError(f"No open spider to crawl: {request}") + self._schedule_request(request, self.spider) + self.slot.nextcall.schedule() # type: ignore[union-attr] - def has_capacity(self): - """Does the engine have capacity to handle more spiders""" - return not bool(self.slot) - - def crawl(self, request, spider): - if spider not in self.open_spiders: - raise RuntimeError(f"Spider {spider.name!r} not opened when crawling: {request}") - self.schedule(request, spider) - self.slot.nextcall.schedule() - - def schedule(self, request, spider): + def _schedule_request(self, request: Request, spider: Spider) -> None: self.signals.send_catch_log(signals.request_scheduled, request=request, spider=spider) - if not self.slot.scheduler.enqueue_request(request): + if not self.slot.scheduler.enqueue_request(request): # type: ignore[union-attr] self.signals.send_catch_log(signals.request_dropped, request=request, spider=spider) - def download(self, request, spider): - d = self._download(request, spider) - d.addBoth(self._downloaded, self.slot, request, spider) - return d + def download(self, request: Request, spider: Optional[Spider] = None) -> Deferred: + """Return a Deferred which fires with a Response as result, only downloader middlewares are applied""" + if spider is None: + spider = self.spider + else: + warnings.warn( + "Passing a 'spider' argument to ExecutionEngine.download is deprecated", + category=ScrapyDeprecationWarning, + stacklevel=2, + ) + if spider is not self.spider: + logger.warning("The spider '%s' does not match the open spider", spider.name) + if spider is None: + raise RuntimeError(f"No open spider to crawl: {request}") + return self._download(request, spider).addBoth(self._downloaded, request, spider) - def _downloaded(self, response, slot, request, spider): - slot.remove_request(request) - return self.download(response, spider) if isinstance(response, Request) else response + def _downloaded( + self, result: Union[Response, Request], request: Request, spider: Spider + ) -> Union[Deferred, Response]: + assert self.slot is not None # typing + self.slot.remove_request(request) + return self.download(result, spider) if isinstance(result, Request) else result - def _download(self, request, spider): - slot = self.slot - slot.add_request(request) + def _download(self, request: Request, spider: Spider) -> Deferred: + assert self.slot is not None # typing - def _on_success(response): - if not isinstance(response, (Response, Request)): - raise TypeError( - "Incorrect type: expected Response or Request, got " - f"{type(response)}: {response!r}" - ) - if isinstance(response, Response): - if response.request is None: - response.request = request - logkws = self.logformatter.crawled(response.request, response, spider) + self.slot.add_request(request) + + def _on_success(result: Union[Response, Request]) -> Union[Response, Request]: + if not isinstance(result, (Response, Request)): + raise TypeError(f"Incorrect type: expected Response or Request, got {type(result)}: {result!r}") + if isinstance(result, Response): + if result.request is None: + result.request = request + logkws = self.logformatter.crawled(result.request, result, spider) if logkws is not None: - logger.log(*logformatter_adapter(logkws), extra={'spider': spider}) + logger.log(*logformatter_adapter(logkws), extra={"spider": spider}) self.signals.send_catch_log( signal=signals.response_received, - response=response, - request=response.request, + response=result, + request=result.request, spider=spider, ) - return response + return result def _on_complete(_): - slot.nextcall.schedule() + self.slot.nextcall.schedule() return _ dwld = self.downloader.fetch(request, spider) @@ -265,58 +289,52 @@ class ExecutionEngine: dwld.addBoth(_on_complete) return dwld - @defer.inlineCallbacks - def open_spider(self, spider, start_requests=(), close_if_idle=True): - if not self.has_capacity(): + @inlineCallbacks + def open_spider(self, spider: Spider, start_requests: Iterable = (), close_if_idle: bool = True): + if self.slot is not None: raise RuntimeError(f"No free spider slot when opening {spider.name!r}") logger.info("Spider opened", extra={'spider': spider}) - nextcall = CallLaterOnce(self._next_request, spider) - scheduler = self.scheduler_cls.from_crawler(self.crawler) + nextcall = CallLaterOnce(self._next_request) + scheduler = create_instance(self.scheduler_cls, settings=None, crawler=self.crawler) start_requests = yield self.scraper.spidermw.process_start_requests(start_requests, spider) - slot = Slot(start_requests, close_if_idle, nextcall, scheduler) - self.slot = slot + self.slot = Slot(start_requests, close_if_idle, nextcall, scheduler) self.spider = spider yield scheduler.open(spider) yield self.scraper.open_spider(spider) self.crawler.stats.open_spider(spider) yield self.signals.send_catch_log_deferred(signals.spider_opened, spider=spider) - slot.nextcall.schedule() - slot.heartbeat.start(5) + self.slot.nextcall.schedule() + self.slot.heartbeat.start(5) - def _spider_idle(self, spider): - """Called when a spider gets idle. This function is called when there - are no remaining pages to download or schedule. It can be called - multiple times. If some extension raises a DontCloseSpider exception - (in the spider_idle signal handler) the spider is not closed until the - next loop and this function is guaranteed to be called (at least) once - again for this spider. + def _spider_idle(self) -> None: """ - res = self.signals.send_catch_log(signals.spider_idle, spider=spider, dont_log=DontCloseSpider) + Called when a spider gets idle, i.e. when there are no remaining requests to download or schedule. + It can be called multiple times. If a handler for the spider_idle signal raises a DontCloseSpider + exception, the spider is not closed until the next loop and this function is guaranteed to be called + (at least) once again. + """ + assert self.spider is not None # typing + res = self.signals.send_catch_log(signals.spider_idle, spider=self.spider, dont_log=DontCloseSpider) if any(isinstance(x, Failure) and isinstance(x.value, DontCloseSpider) for _, x in res): - return + return None + if self.spider_is_idle(): + self.close_spider(self.spider, reason='finished') - if self.spider_is_idle(spider): - self.close_spider(spider, reason='finished') - - def close_spider(self, spider, reason='cancelled'): + def close_spider(self, spider: Spider, reason: str = "cancelled") -> Deferred: """Close (cancel) spider and clear all its outstanding requests""" + if self.slot is None: + raise RuntimeError("Engine slot not assigned") - slot = self.slot - if slot.closing: - return slot.closing - logger.info("Closing spider (%(reason)s)", - {'reason': reason}, - extra={'spider': spider}) + if self.slot.closing is not None: + return self.slot.closing - dfd = slot.close() + logger.info("Closing spider (%(reason)s)", {'reason': reason}, extra={'spider': spider}) - def log_failure(msg): - def errback(failure): - logger.error( - msg, - exc_info=failure_to_exc_info(failure), - extra={'spider': spider} - ) + dfd = self.slot.close() + + def log_failure(msg: str) -> Callable: + def errback(failure: Failure) -> None: + logger.error(msg, exc_info=failure_to_exc_info(failure), extra={'spider': spider}) return errback dfd.addBoth(lambda _: self.downloader.close()) @@ -325,19 +343,18 @@ class ExecutionEngine: dfd.addBoth(lambda _: self.scraper.close_spider(spider)) dfd.addErrback(log_failure('Scraper close failure')) - dfd.addBoth(lambda _: slot.scheduler.close(reason)) + dfd.addBoth(lambda _: self.slot.scheduler.close(reason)) dfd.addErrback(log_failure('Scheduler close failure')) dfd.addBoth(lambda _: self.signals.send_catch_log_deferred( - signal=signals.spider_closed, spider=spider, reason=reason)) + signal=signals.spider_closed, spider=spider, reason=reason, + )) dfd.addErrback(log_failure('Error while sending spider_close signal')) dfd.addBoth(lambda _: self.crawler.stats.close_spider(spider, reason=reason)) dfd.addErrback(log_failure('Stats close failure')) - dfd.addBoth(lambda _: logger.info("Spider closed (%(reason)s)", - {'reason': reason}, - extra={'spider': spider})) + dfd.addBoth(lambda _: logger.info("Spider closed (%(reason)s)", {'reason': reason}, extra={'spider': spider})) dfd.addBoth(lambda _: setattr(self, 'slot', None)) dfd.addErrback(log_failure('Error while unassigning slot')) @@ -349,12 +366,26 @@ class ExecutionEngine: return dfd - def _close_all_spiders(self): - dfds = [self.close_spider(s, reason='shutdown') for s in self.open_spiders] - dlist = defer.DeferredList(dfds) - return dlist + @property + def open_spiders(self) -> list: + warnings.warn( + "ExecutionEngine.open_spiders is deprecated, please use ExecutionEngine.spider instead", + category=ScrapyDeprecationWarning, + stacklevel=2, + ) + return [self.spider] if self.spider is not None else [] - @defer.inlineCallbacks - def _finish_stopping_engine(self): - yield self.signals.send_catch_log_deferred(signal=signals.engine_stopped) - self._closewait.callback(None) + def has_capacity(self) -> bool: + warnings.warn("ExecutionEngine.has_capacity is deprecated", ScrapyDeprecationWarning, stacklevel=2) + return not bool(self.slot) + + def schedule(self, request: Request, spider: Spider) -> None: + warnings.warn( + "ExecutionEngine.schedule is deprecated, please use " + "ExecutionEngine.crawl or ExecutionEngine.download instead", + category=ScrapyDeprecationWarning, + stacklevel=2, + ) + if self.slot is None: + raise RuntimeError("Engine slot not assigned") + self._schedule_request(request, spider) diff --git a/scrapy/core/scraper.py b/scrapy/core/scraper.py index 96aa53686..d6d6f64f9 100644 --- a/scrapy/core/scraper.py +++ b/scrapy/core/scraper.py @@ -200,7 +200,7 @@ class Scraper: """ assert self.slot is not None # typing if isinstance(output, Request): - self.crawler.engine.crawl(request=output, spider=spider) + self.crawler.engine.crawl(request=output) elif is_item(output): self.slot.itemproc_size += 1 dfd = self.itemproc.process_item(output, spider) diff --git a/scrapy/downloadermiddlewares/robotstxt.py b/scrapy/downloadermiddlewares/robotstxt.py index d6da55535..e66bf177e 100644 --- a/scrapy/downloadermiddlewares/robotstxt.py +++ b/scrapy/downloadermiddlewares/robotstxt.py @@ -67,7 +67,7 @@ class RobotsTxtMiddleware: priority=self.DOWNLOAD_PRIORITY, meta={'dont_obey_robotstxt': True} ) - dfd = self.crawler.engine.download(robotsreq, spider) + dfd = self.crawler.engine.download(robotsreq) dfd.addCallback(self._parse_robots, netloc, spider) dfd.addErrback(self._logerror, robotsreq, spider) dfd.addErrback(self._robots_error, netloc) diff --git a/scrapy/extensions/memusage.py b/scrapy/extensions/memusage.py index 274cbdbfe..9de119a10 100644 --- a/scrapy/extensions/memusage.py +++ b/scrapy/extensions/memusage.py @@ -88,10 +88,8 @@ class MemoryUsage: self._send_report(self.notify_mails, subj) self.crawler.stats.set_value('memusage/limit_notified', 1) - open_spiders = self.crawler.engine.open_spiders - if open_spiders: - for spider in open_spiders: - self.crawler.engine.close_spider(spider, 'memusage_exceeded') + if self.crawler.engine.spider is not None: + self.crawler.engine.close_spider(self.crawler.engine.spider, 'memusage_exceeded') else: self.crawler.stop() diff --git a/scrapy/pipelines/media.py b/scrapy/pipelines/media.py index 0c2ee6856..d1bccf323 100644 --- a/scrapy/pipelines/media.py +++ b/scrapy/pipelines/media.py @@ -173,7 +173,7 @@ class MediaPipeline: errback=self.media_failed, errbackArgs=(request, info)) else: self._modify_media_request(request) - dfd = self.crawler.engine.download(request, info.spider) + dfd = self.crawler.engine.download(request) dfd.addCallbacks( callback=self.media_downloaded, callbackArgs=(request, info), callbackKeywords={'item': item}, errback=self.media_failed, errbackArgs=(request, info)) diff --git a/scrapy/shell.py b/scrapy/shell.py index c370ccaff..f2dff2ae3 100644 --- a/scrapy/shell.py +++ b/scrapy/shell.py @@ -79,7 +79,7 @@ class Shell: spider = self._open_spider(request, spider) d = _request_deferred(request) d.addCallback(lambda x: (x, spider)) - self.crawler.engine.crawl(request, spider) + self.crawler.engine.crawl(request) return d def _open_spider(self, request, spider): diff --git a/scrapy/utils/engine.py b/scrapy/utils/engine.py index 0c1cee1a0..8e3ec2c37 100644 --- a/scrapy/utils/engine.py +++ b/scrapy/utils/engine.py @@ -8,11 +8,10 @@ def get_engine_status(engine): """Return a report of the current engine status""" tests = [ "time()-engine.start_time", - "engine.has_capacity()", "len(engine.downloader.active)", "engine.scraper.is_idle()", "engine.spider.name", - "engine.spider_is_idle(engine.spider)", + "engine.spider_is_idle()", "engine.slot.closing", "len(engine.slot.inprogress)", "len(engine.slot.scheduler.dqs or [])", diff --git a/tests/test_downloadermiddleware_robotstxt.py b/tests/test_downloadermiddleware_robotstxt.py index 858138f81..1460d88eb 100644 --- a/tests/test_downloadermiddleware_robotstxt.py +++ b/tests/test_downloadermiddleware_robotstxt.py @@ -42,7 +42,7 @@ Disallow: /some/randome/page.html """.encode('utf-8') response = TextResponse('http://site.local/robots.txt', body=ROBOTS) - def return_response(request, spider): + def return_response(request): deferred = Deferred() reactor.callFromThread(deferred.callback, response) return deferred @@ -79,7 +79,7 @@ Disallow: /some/randome/page.html crawler.settings.set('ROBOTSTXT_OBEY', True) response = Response('http://site.local/robots.txt', body=b'GIF89a\xd3\x00\xfe\x00\xa2') - def return_response(request, spider): + def return_response(request): deferred = Deferred() reactor.callFromThread(deferred.callback, response) return deferred @@ -102,7 +102,7 @@ Disallow: /some/randome/page.html crawler.settings.set('ROBOTSTXT_OBEY', True) response = Response('http://site.local/robots.txt') - def return_response(request, spider): + def return_response(request): deferred = Deferred() reactor.callFromThread(deferred.callback, response) return deferred @@ -122,7 +122,7 @@ Disallow: /some/randome/page.html self.crawler.settings.set('ROBOTSTXT_OBEY', True) err = error.DNSLookupError('Robotstxt address not found') - def return_failure(request, spider): + def return_failure(request): deferred = Deferred() reactor.callFromThread(deferred.errback, failure.Failure(err)) return deferred @@ -138,7 +138,7 @@ Disallow: /some/randome/page.html self.crawler.settings.set('ROBOTSTXT_OBEY', True) err = error.DNSLookupError('Robotstxt address not found') - def immediate_failure(request, spider): + def immediate_failure(request): deferred = Deferred() deferred.errback(failure.Failure(err)) return deferred @@ -150,7 +150,7 @@ Disallow: /some/randome/page.html def test_ignore_robotstxt_request(self): self.crawler.settings.set('ROBOTSTXT_OBEY', True) - def ignore_request(request, spider): + def ignore_request(request): deferred = Deferred() reactor.callFromThread(deferred.errback, failure.Failure(IgnoreRequest())) return deferred diff --git a/tests/test_engine.py b/tests/test_engine.py index b2d1d83c7..c200ded90 100644 --- a/tests/test_engine.py +++ b/tests/test_engine.py @@ -13,6 +13,7 @@ module with the ``runserver`` argument:: import os import re import sys +import warnings from collections import defaultdict from urllib.parse import urlparse @@ -25,6 +26,7 @@ from twisted.web import server, static, util from scrapy import signals from scrapy.core.engine import ExecutionEngine +from scrapy.exceptions import ScrapyDeprecationWarning from scrapy.http import Request from scrapy.item import Item, Field from scrapy.linkextractors import LinkExtractor @@ -382,22 +384,104 @@ class EngineTest(unittest.TestCase): yield e.close() @defer.inlineCallbacks - def test_close_spiders_downloader(self): - e = ExecutionEngine(get_crawler(TestSpider), lambda _: None) - yield e.open_spider(TestSpider(), []) - self.assertEqual(len(e.open_spiders), 1) - yield e.close() - self.assertEqual(len(e.open_spiders), 0) - - @defer.inlineCallbacks - def test_close_engine_spiders_downloader(self): + def test_start_already_running_exception(self): e = ExecutionEngine(get_crawler(TestSpider), lambda _: None) yield e.open_spider(TestSpider(), []) e.start() - self.assertTrue(e.running) - yield e.close() - self.assertFalse(e.running) - self.assertEqual(len(e.open_spiders), 0) + yield self.assertFailure(e.start(), RuntimeError).addBoth( + lambda exc: self.assertEqual(str(exc), "Engine already running") + ) + yield e.stop() + + @defer.inlineCallbacks + def test_close_spiders_downloader(self): + with warnings.catch_warnings(record=True) as warning_list: + e = ExecutionEngine(get_crawler(TestSpider), lambda _: None) + yield e.open_spider(TestSpider(), []) + self.assertEqual(len(e.open_spiders), 1) + yield e.close() + self.assertEqual(len(e.open_spiders), 0) + self.assertEqual(warning_list[0].category, ScrapyDeprecationWarning) + self.assertEqual( + str(warning_list[0].message), + "ExecutionEngine.open_spiders is deprecated, please use ExecutionEngine.spider instead", + ) + + @defer.inlineCallbacks + def test_close_engine_spiders_downloader(self): + with warnings.catch_warnings(record=True) as warning_list: + e = ExecutionEngine(get_crawler(TestSpider), lambda _: None) + yield e.open_spider(TestSpider(), []) + e.start() + self.assertTrue(e.running) + yield e.close() + self.assertFalse(e.running) + self.assertEqual(len(e.open_spiders), 0) + self.assertEqual(warning_list[0].category, ScrapyDeprecationWarning) + self.assertEqual( + str(warning_list[0].message), + "ExecutionEngine.open_spiders is deprecated, please use ExecutionEngine.spider instead", + ) + + @defer.inlineCallbacks + def test_crawl_deprecated_spider_arg(self): + with warnings.catch_warnings(record=True) as warning_list: + e = ExecutionEngine(get_crawler(TestSpider), lambda _: None) + spider = TestSpider() + yield e.open_spider(spider, []) + e.start() + e.crawl(Request("data:,"), spider) + yield e.close() + self.assertEqual(warning_list[0].category, ScrapyDeprecationWarning) + self.assertEqual( + str(warning_list[0].message), + "Passing a 'spider' argument to ExecutionEngine.crawl is deprecated", + ) + + @defer.inlineCallbacks + def test_download_deprecated_spider_arg(self): + with warnings.catch_warnings(record=True) as warning_list: + e = ExecutionEngine(get_crawler(TestSpider), lambda _: None) + spider = TestSpider() + yield e.open_spider(spider, []) + e.start() + e.download(Request("data:,"), spider) + yield e.close() + self.assertEqual(warning_list[0].category, ScrapyDeprecationWarning) + self.assertEqual( + str(warning_list[0].message), + "Passing a 'spider' argument to ExecutionEngine.download is deprecated", + ) + + @defer.inlineCallbacks + def test_deprecated_schedule(self): + with warnings.catch_warnings(record=True) as warning_list: + e = ExecutionEngine(get_crawler(TestSpider), lambda _: None) + spider = TestSpider() + yield e.open_spider(spider, []) + e.start() + e.schedule(Request("data:,"), spider) + yield e.close() + self.assertEqual(warning_list[0].category, ScrapyDeprecationWarning) + self.assertEqual( + str(warning_list[0].message), + "ExecutionEngine.schedule is deprecated, please use " + "ExecutionEngine.crawl or ExecutionEngine.download instead", + ) + + @defer.inlineCallbacks + def test_deprecated_has_capacity(self): + with warnings.catch_warnings(record=True) as warning_list: + e = ExecutionEngine(get_crawler(TestSpider), lambda _: None) + self.assertTrue(e.has_capacity()) + spider = TestSpider() + yield e.open_spider(spider, []) + self.assertFalse(e.has_capacity()) + e.start() + yield e.close() + self.assertTrue(e.has_capacity()) + self.assertEqual(warning_list[0].category, ScrapyDeprecationWarning) + self.assertEqual(str(warning_list[0].message), "ExecutionEngine.has_capacity is deprecated") if __name__ == "__main__": From e3f81d8d5f17515b6eba135ac0db7e270ff0a9f0 Mon Sep 17 00:00:00 2001 From: Eugenio Lacuesta <1731933+elacuesta@users.noreply.github.com> Date: Tue, 20 Apr 2021 11:46:43 -0300 Subject: [PATCH 05/16] Engine: remove unnecessary parameter (#5106) --- scrapy/core/engine.py | 10 ++++++---- 1 file changed, 6 insertions(+), 4 deletions(-) diff --git a/scrapy/core/engine.py b/scrapy/core/engine.py index edfac87c6..7a09bafa1 100644 --- a/scrapy/core/engine.py +++ b/scrapy/core/engine.py @@ -161,7 +161,7 @@ class ExecutionEngine: return None d = self._download(request, self.spider) - d.addBoth(self._handle_downloader_output, request, self.spider) + d.addBoth(self._handle_downloader_output, request) d.addErrback(lambda f: logger.info('Error while handling downloader output', exc_info=failure_to_exc_info(f), extra={'spider': self.spider})) @@ -176,8 +176,10 @@ class ExecutionEngine: return d def _handle_downloader_output( - self, result: Union[Request, Response, Failure], request: Request, spider: Spider + self, result: Union[Request, Response, Failure], request: Request ) -> Optional[Deferred]: + assert self.spider is not None # typing + if not isinstance(result, (Request, Response, Failure)): raise TypeError(f"Incorrect type: expected Request, Response or Failure, got {type(result)}: {result!r}") @@ -186,12 +188,12 @@ class ExecutionEngine: self.crawl(result) return None - d = self.scraper.enqueue_scrape(result, request, spider) + d = self.scraper.enqueue_scrape(result, request, self.spider) d.addErrback( lambda f: logger.error( "Error while enqueuing downloader output", exc_info=failure_to_exc_info(f), - extra={'spider': spider}, + extra={'spider': self.spider}, ) ) return d From e779ed7d93beec36f565d33a1cc8d3e8fe6068d7 Mon Sep 17 00:00:00 2001 From: Eugenio Lacuesta <1731933+elacuesta@users.noreply.github.com> Date: Tue, 20 Apr 2021 16:39:07 -0300 Subject: [PATCH 06/16] Dupefilter type hints (#5108) --- scrapy/dupefilters.py | 41 +++++++++++++++++++++++++++-------------- scrapy/utils/request.py | 2 +- 2 files changed, 28 insertions(+), 15 deletions(-) diff --git a/scrapy/dupefilters.py b/scrapy/dupefilters.py index ac5478e7c..292c68099 100644 --- a/scrapy/dupefilters.py +++ b/scrapy/dupefilters.py @@ -1,35 +1,47 @@ -import os import logging +import os +from typing import Optional, Set, Type, TypeVar +from twisted.internet.defer import Deferred + +from scrapy.http.request import Request +from scrapy.settings import BaseSettings +from scrapy.spiders import Spider from scrapy.utils.job import job_dir from scrapy.utils.request import referer_str, request_fingerprint -class BaseDupeFilter: +BaseDupeFilterTV = TypeVar("BaseDupeFilterTV", bound="BaseDupeFilter") + +class BaseDupeFilter: @classmethod - def from_settings(cls, settings): + def from_settings(cls: Type[BaseDupeFilterTV], settings: BaseSettings) -> BaseDupeFilterTV: return cls() - def request_seen(self, request): + def request_seen(self, request: Request) -> bool: return False - def open(self): # can return deferred + def open(self) -> Optional[Deferred]: pass - def close(self, reason): # can return a deferred + def close(self, reason: str) -> Optional[Deferred]: pass - def log(self, request, spider): # log that a request has been filtered + def log(self, request: Request, spider: Spider) -> None: + """Log that a request has been filtered""" pass +RFPDupeFilterTV = TypeVar("RFPDupeFilterTV", bound="RFPDupeFilter") + + class RFPDupeFilter(BaseDupeFilter): """Request Fingerprint duplicates filter""" - def __init__(self, path=None, debug=False): + def __init__(self, path: Optional[str] = None, debug: bool = False) -> None: self.file = None - self.fingerprints = set() + self.fingerprints: Set[str] = set() self.logdupes = True self.debug = debug self.logger = logging.getLogger(__name__) @@ -39,26 +51,27 @@ class RFPDupeFilter(BaseDupeFilter): self.fingerprints.update(x.rstrip() for x in self.file) @classmethod - def from_settings(cls, settings): + def from_settings(cls: Type[RFPDupeFilterTV], settings: BaseSettings) -> RFPDupeFilterTV: debug = settings.getbool('DUPEFILTER_DEBUG') return cls(job_dir(settings), debug) - def request_seen(self, request): + def request_seen(self, request: Request) -> bool: fp = self.request_fingerprint(request) if fp in self.fingerprints: return True self.fingerprints.add(fp) if self.file: self.file.write(fp + '\n') + return False - def request_fingerprint(self, request): + def request_fingerprint(self, request: Request) -> str: return request_fingerprint(request) - def close(self, reason): + def close(self, reason: str) -> None: if self.file: self.file.close() - def log(self, request, spider): + def log(self, request: Request, spider: Spider) -> None: if self.debug: msg = "Filtered duplicate request: %(request)s (referer: %(referer)s)" args = {'request': request, 'referer': referer_str(request)} diff --git a/scrapy/utils/request.py b/scrapy/utils/request.py index 66736b42f..541368423 100644 --- a/scrapy/utils/request.py +++ b/scrapy/utils/request.py @@ -24,7 +24,7 @@ def request_fingerprint( request: Request, include_headers: Optional[Iterable[Union[bytes, str]]] = None, keep_fragments: bool = False, -): +) -> str: """ Return the request fingerprint. From 68379197986ae3deb81a545b5fd6920ea3347094 Mon Sep 17 00:00:00 2001 From: Eugenio Lacuesta <1731933+elacuesta@users.noreply.github.com> Date: Mon, 26 Apr 2021 14:55:02 -0300 Subject: [PATCH 07/16] Add peek method to queues (#5112) --- pylintrc | 1 + scrapy/pqueues.py | 59 +++++++--- scrapy/squeues.py | 59 +++++++--- tests/test_pqueues.py | 144 +++++++++++++++++++++++ tests/test_squeues_request.py | 214 ++++++++++++++++++++++++++++++++++ 5 files changed, 447 insertions(+), 30 deletions(-) create mode 100644 tests/test_pqueues.py create mode 100644 tests/test_squeues_request.py diff --git a/pylintrc b/pylintrc index 5b6b9fab0..972bf99de 100644 --- a/pylintrc +++ b/pylintrc @@ -24,6 +24,7 @@ disable=abstract-method, consider-using-in, consider-using-set-comprehension, consider-using-sys-exit, + consider-using-with, cyclic-import, dangerous-default-value, deprecated-method, diff --git a/scrapy/pqueues.py b/scrapy/pqueues.py index a9aa6c649..b4b63e7c7 100644 --- a/scrapy/pqueues.py +++ b/scrapy/pqueues.py @@ -3,6 +3,7 @@ import logging from scrapy.utils.misc import create_instance + logger = logging.getLogger(__name__) @@ -17,8 +18,7 @@ def _path_safe(text): >>> _path_safe('some@symbol?').startswith('some_symbol_') True """ - pathable_slot = "".join([c if c.isalnum() or c in '-._' else '_' - for c in text]) + pathable_slot = "".join([c if c.isalnum() or c in '-._' else '_' for c in text]) # as we replace some letters we can get collision for different slots # add we add unique part unique_slot = hashlib.md5(text.encode('utf8')).hexdigest() @@ -35,6 +35,9 @@ class ScrapyPriorityQueue: * close() * __len__() + Optionally, the queue could provide a ``peek`` method, that should return the + next object to be returned by ``pop``, but without removing it from the queue. + ``__init__`` method of ScrapyPriorityQueue receives a downstream_queue_cls argument, which is a class used to instantiate a new (internal) queue when a new priority is allocated. @@ -70,10 +73,12 @@ class ScrapyPriorityQueue: self.curprio = min(startprios) def qfactory(self, key): - return create_instance(self.downstream_queue_cls, - None, - self.crawler, - self.key + '/' + str(key)) + return create_instance( + self.downstream_queue_cls, + None, + self.crawler, + self.key + '/' + str(key), + ) def priority(self, request): return -request.priority @@ -99,6 +104,18 @@ class ScrapyPriorityQueue: self.curprio = min(prios) if prios else None return m + def peek(self): + """Returns the next object to be returned by :meth:`pop`, + but without removing it from the queue. + + Raises :exc:`NotImplementedError` if the underlying queue class does + not implement a ``peek`` method, which is optional for queues. + """ + if self.curprio is None: + return None + queue = self.queues[self.curprio] + return queue.peek() + def close(self): active = [] for p, q in self.queues.items(): @@ -116,8 +133,7 @@ class DownloaderInterface: self.downloader = crawler.engine.downloader def stats(self, possible_slots): - return [(self._active_downloads(slot), slot) - for slot in possible_slots] + return [(self._active_downloads(slot), slot) for slot in possible_slots] def get_slot_key(self, request): return self.downloader._get_slot_key(request, None) @@ -162,10 +178,12 @@ class DownloaderAwarePriorityQueue: self.pqueues[slot] = self.pqfactory(slot, startprios) def pqfactory(self, slot, startprios=()): - return ScrapyPriorityQueue(self.crawler, - self.downstream_queue_cls, - self.key + '/' + _path_safe(slot), - startprios) + return ScrapyPriorityQueue( + self.crawler, + self.downstream_queue_cls, + self.key + '/' + _path_safe(slot), + startprios, + ) def pop(self): stats = self._downloader_interface.stats(self.pqueues) @@ -187,9 +205,22 @@ class DownloaderAwarePriorityQueue: queue = self.pqueues[slot] queue.push(request) + def peek(self): + """Returns the next object to be returned by :meth:`pop`, + but without removing it from the queue. + + Raises :exc:`NotImplementedError` if the underlying queue class does + not implement a ``peek`` method, which is optional for queues. + """ + stats = self._downloader_interface.stats(self.pqueues) + if not stats: + return None + slot = min(stats)[1] + queue = self.pqueues[slot] + return queue.peek() + def close(self): - active = {slot: queue.close() - for slot, queue in self.pqueues.items()} + active = {slot: queue.close() for slot, queue in self.pqueues.items()} self.pqueues.clear() return active diff --git a/scrapy/squeues.py b/scrapy/squeues.py index 77ffda6f7..44898ba08 100644 --- a/scrapy/squeues.py +++ b/scrapy/squeues.py @@ -19,7 +19,6 @@ def _with_mkdir(queue_class): dirname = os.path.dirname(path) if not os.path.exists(dirname): os.makedirs(dirname, exist_ok=True) - super().__init__(path, *args, **kwargs) return DirectoriesCreated @@ -38,6 +37,20 @@ def _serializable_queue(queue_class, serialize, deserialize): if s: return deserialize(s) + def peek(self): + """Returns the next object to be returned by :meth:`pop`, + but without removing it from the queue. + + Raises :exc:`NotImplementedError` if the underlying queue class does + not implement a ``peek`` method, which is optional for queues. + """ + try: + s = super().peek() + except AttributeError as ex: + raise NotImplementedError("The underlying queue class does not implement 'peek'") from ex + if s: + return deserialize(s) + return SerializableQueue @@ -59,12 +72,21 @@ def _scrapy_serialization_queue(queue_class): def pop(self): request = super().pop() - if not request: return None + return request_from_dict(request, self.spider) - request = request_from_dict(request, self.spider) - return request + def peek(self): + """Returns the next object to be returned by :meth:`pop`, + but without removing it from the queue. + + Raises :exc:`NotImplementedError` if the underlying queue class does + not implement a ``peek`` method, which is optional for queues. + """ + request = super().peek() + if not request: + return None + return request_from_dict(request, self.spider) return ScrapyRequestQueue @@ -76,6 +98,19 @@ def _scrapy_non_serialization_queue(queue_class): def from_crawler(cls, crawler, *args, **kwargs): return cls() + def peek(self): + """Returns the next object to be returned by :meth:`pop`, + but without removing it from the queue. + + Raises :exc:`NotImplementedError` if the underlying queue class does + not implement a ``peek`` method, which is optional for queues. + """ + try: + s = super().peek() + except AttributeError as ex: + raise NotImplementedError("The underlying queue class does not implement 'peek'") from ex + return s + return ScrapyRequestQueue @@ -109,17 +144,9 @@ MarshalLifoDiskQueueNonRequest = _serializable_queue( marshal.loads ) -PickleFifoDiskQueue = _scrapy_serialization_queue( - PickleFifoDiskQueueNonRequest -) -PickleLifoDiskQueue = _scrapy_serialization_queue( - PickleLifoDiskQueueNonRequest -) -MarshalFifoDiskQueue = _scrapy_serialization_queue( - MarshalFifoDiskQueueNonRequest -) -MarshalLifoDiskQueue = _scrapy_serialization_queue( - MarshalLifoDiskQueueNonRequest -) +PickleFifoDiskQueue = _scrapy_serialization_queue(PickleFifoDiskQueueNonRequest) +PickleLifoDiskQueue = _scrapy_serialization_queue(PickleLifoDiskQueueNonRequest) +MarshalFifoDiskQueue = _scrapy_serialization_queue(MarshalFifoDiskQueueNonRequest) +MarshalLifoDiskQueue = _scrapy_serialization_queue(MarshalLifoDiskQueueNonRequest) FifoMemoryQueue = _scrapy_non_serialization_queue(queue.FifoMemoryQueue) LifoMemoryQueue = _scrapy_non_serialization_queue(queue.LifoMemoryQueue) diff --git a/tests/test_pqueues.py b/tests/test_pqueues.py new file mode 100644 index 000000000..ec55033d1 --- /dev/null +++ b/tests/test_pqueues.py @@ -0,0 +1,144 @@ +import tempfile +import unittest + +import queuelib + +from scrapy.http.request import Request +from scrapy.pqueues import ScrapyPriorityQueue, DownloaderAwarePriorityQueue +from scrapy.spiders import Spider +from scrapy.squeues import FifoMemoryQueue +from scrapy.utils.test import get_crawler + +from tests.test_scheduler import MockDownloader, MockEngine + + +class PriorityQueueTest(unittest.TestCase): + def setUp(self): + self.crawler = get_crawler(Spider) + self.spider = self.crawler._create_spider("foo") + + def test_queue_push_pop_one(self): + temp_dir = tempfile.mkdtemp() + queue = ScrapyPriorityQueue.from_crawler(self.crawler, FifoMemoryQueue, temp_dir) + self.assertIsNone(queue.pop()) + self.assertEqual(len(queue), 0) + req1 = Request("https://example.org/1", priority=1) + queue.push(req1) + self.assertEqual(len(queue), 1) + dequeued = queue.pop() + self.assertEqual(len(queue), 0) + self.assertEqual(dequeued.url, req1.url) + self.assertEqual(dequeued.priority, req1.priority) + self.assertEqual(queue.close(), []) + + def test_no_peek_raises(self): + if hasattr(queuelib.queue.FifoMemoryQueue, "peek"): + raise unittest.SkipTest("queuelib.queue.FifoMemoryQueue.peek is defined") + temp_dir = tempfile.mkdtemp() + queue = ScrapyPriorityQueue.from_crawler(self.crawler, FifoMemoryQueue, temp_dir) + queue.push(Request("https://example.org")) + with self.assertRaises(NotImplementedError, msg="The underlying queue class does not implement 'peek'"): + queue.peek() + queue.close() + + def test_peek(self): + if not hasattr(queuelib.queue.FifoMemoryQueue, "peek"): + raise unittest.SkipTest("queuelib.queue.FifoMemoryQueue.peek is undefined") + temp_dir = tempfile.mkdtemp() + queue = ScrapyPriorityQueue.from_crawler(self.crawler, FifoMemoryQueue, temp_dir) + self.assertEqual(len(queue), 0) + self.assertIsNone(queue.peek()) + req1 = Request("https://example.org/1") + req2 = Request("https://example.org/2") + req3 = Request("https://example.org/3") + queue.push(req1) + queue.push(req2) + queue.push(req3) + self.assertEqual(len(queue), 3) + self.assertEqual(queue.peek().url, req1.url) + self.assertEqual(queue.pop().url, req1.url) + self.assertEqual(len(queue), 2) + self.assertEqual(queue.peek().url, req2.url) + self.assertEqual(queue.pop().url, req2.url) + self.assertEqual(len(queue), 1) + self.assertEqual(queue.peek().url, req3.url) + self.assertEqual(queue.pop().url, req3.url) + self.assertEqual(queue.close(), []) + + def test_queue_push_pop_priorities(self): + temp_dir = tempfile.mkdtemp() + queue = ScrapyPriorityQueue.from_crawler(self.crawler, FifoMemoryQueue, temp_dir, [-1, -2, -3]) + self.assertIsNone(queue.pop()) + self.assertEqual(len(queue), 0) + req1 = Request("https://example.org/1", priority=1) + req2 = Request("https://example.org/2", priority=2) + req3 = Request("https://example.org/3", priority=3) + queue.push(req1) + queue.push(req2) + queue.push(req3) + self.assertEqual(len(queue), 3) + dequeued = queue.pop() + self.assertEqual(len(queue), 2) + self.assertEqual(dequeued.url, req3.url) + self.assertEqual(dequeued.priority, req3.priority) + self.assertEqual(queue.close(), [-1, -2]) + + +class DownloaderAwarePriorityQueueTest(unittest.TestCase): + def setUp(self): + crawler = get_crawler(Spider) + crawler.engine = MockEngine(downloader=MockDownloader()) + self.queue = DownloaderAwarePriorityQueue.from_crawler( + crawler=crawler, + downstream_queue_cls=FifoMemoryQueue, + key="foo/bar", + ) + + def tearDown(self): + self.queue.close() + + def test_push_pop(self): + self.assertEqual(len(self.queue), 0) + self.assertIsNone(self.queue.pop()) + req1 = Request("http://www.example.com/1") + req2 = Request("http://www.example.com/2") + req3 = Request("http://www.example.com/3") + self.queue.push(req1) + self.queue.push(req2) + self.queue.push(req3) + self.assertEqual(len(self.queue), 3) + self.assertEqual(self.queue.pop().url, req1.url) + self.assertEqual(len(self.queue), 2) + self.assertEqual(self.queue.pop().url, req2.url) + self.assertEqual(len(self.queue), 1) + self.assertEqual(self.queue.pop().url, req3.url) + self.assertEqual(len(self.queue), 0) + self.assertIsNone(self.queue.pop()) + + def test_no_peek_raises(self): + if hasattr(queuelib.queue.FifoMemoryQueue, "peek"): + raise unittest.SkipTest("queuelib.queue.FifoMemoryQueue.peek is defined") + self.queue.push(Request("https://example.org")) + with self.assertRaises(NotImplementedError, msg="The underlying queue class does not implement 'peek'"): + self.queue.peek() + + def test_peek(self): + if not hasattr(queuelib.queue.FifoMemoryQueue, "peek"): + raise unittest.SkipTest("queuelib.queue.FifoMemoryQueue.peek is undefined") + self.assertEqual(len(self.queue), 0) + req1 = Request("https://example.org/1") + req2 = Request("https://example.org/2") + req3 = Request("https://example.org/3") + self.queue.push(req1) + self.queue.push(req2) + self.queue.push(req3) + self.assertEqual(len(self.queue), 3) + self.assertEqual(self.queue.peek().url, req1.url) + self.assertEqual(self.queue.pop().url, req1.url) + self.assertEqual(len(self.queue), 2) + self.assertEqual(self.queue.peek().url, req2.url) + self.assertEqual(self.queue.pop().url, req2.url) + self.assertEqual(len(self.queue), 1) + self.assertEqual(self.queue.peek().url, req3.url) + self.assertEqual(self.queue.pop().url, req3.url) + self.assertIsNone(self.queue.peek()) diff --git a/tests/test_squeues_request.py b/tests/test_squeues_request.py new file mode 100644 index 000000000..c5fcc1853 --- /dev/null +++ b/tests/test_squeues_request.py @@ -0,0 +1,214 @@ +import shutil +import tempfile +import unittest + +import queuelib + +from scrapy.squeues import ( + PickleFifoDiskQueue, + PickleLifoDiskQueue, + MarshalFifoDiskQueue, + MarshalLifoDiskQueue, + FifoMemoryQueue, + LifoMemoryQueue, +) +from scrapy.http import Request +from scrapy.spiders import Spider +from scrapy.utils.test import get_crawler + + +""" +Queues that handle requests +""" + + +class BaseQueueTestCase(unittest.TestCase): + def setUp(self): + self.tmpdir = tempfile.mkdtemp(prefix="scrapy-queue-tests-") + self.qpath = self.tempfilename() + self.qdir = self.mkdtemp() + self.crawler = get_crawler(Spider) + + def tearDown(self): + shutil.rmtree(self.tmpdir) + + def tempfilename(self): + with tempfile.NamedTemporaryFile(dir=self.tmpdir) as nf: + return nf.name + + def mkdtemp(self): + return tempfile.mkdtemp(dir=self.tmpdir) + + +class RequestQueueTestMixin: + def queue(self): + raise NotImplementedError() + + def test_one_element_with_peek(self): + if not hasattr(queuelib.queue.FifoMemoryQueue, "peek"): + raise unittest.SkipTest("The queuelib queues do not define peek") + q = self.queue() + self.assertEqual(len(q), 0) + self.assertIsNone(q.peek()) + self.assertIsNone(q.pop()) + req = Request("http://www.example.com") + q.push(req) + self.assertEqual(len(q), 1) + self.assertEqual(q.peek().url, req.url) + self.assertEqual(q.pop().url, req.url) + self.assertEqual(len(q), 0) + self.assertIsNone(q.peek()) + self.assertIsNone(q.pop()) + q.close() + + def test_one_element_without_peek(self): + if hasattr(queuelib.queue.FifoMemoryQueue, "peek"): + raise unittest.SkipTest("The queuelib queues define peek") + q = self.queue() + self.assertEqual(len(q), 0) + self.assertIsNone(q.pop()) + req = Request("http://www.example.com") + q.push(req) + self.assertEqual(len(q), 1) + with self.assertRaises(NotImplementedError, msg="The underlying queue class does not implement 'peek'"): + q.peek() + self.assertEqual(q.pop().url, req.url) + self.assertEqual(len(q), 0) + self.assertIsNone(q.pop()) + q.close() + + +class FifoQueueMixin(RequestQueueTestMixin): + def test_fifo_with_peek(self): + if not hasattr(queuelib.queue.FifoMemoryQueue, "peek"): + raise unittest.SkipTest("The queuelib queues do not define peek") + q = self.queue() + self.assertEqual(len(q), 0) + self.assertIsNone(q.peek()) + self.assertIsNone(q.pop()) + req1 = Request("http://www.example.com/1") + req2 = Request("http://www.example.com/2") + req3 = Request("http://www.example.com/3") + q.push(req1) + q.push(req2) + q.push(req3) + self.assertEqual(len(q), 3) + self.assertEqual(q.peek().url, req1.url) + self.assertEqual(q.pop().url, req1.url) + self.assertEqual(len(q), 2) + self.assertEqual(q.peek().url, req2.url) + self.assertEqual(q.pop().url, req2.url) + self.assertEqual(len(q), 1) + self.assertEqual(q.peek().url, req3.url) + self.assertEqual(q.pop().url, req3.url) + self.assertEqual(len(q), 0) + self.assertIsNone(q.peek()) + self.assertIsNone(q.pop()) + q.close() + + def test_fifo_without_peek(self): + if hasattr(queuelib.queue.FifoMemoryQueue, "peek"): + raise unittest.SkipTest("The queuelib queues do not define peek") + q = self.queue() + self.assertEqual(len(q), 0) + self.assertIsNone(q.pop()) + req1 = Request("http://www.example.com/1") + req2 = Request("http://www.example.com/2") + req3 = Request("http://www.example.com/3") + q.push(req1) + q.push(req2) + q.push(req3) + with self.assertRaises(NotImplementedError, msg="The underlying queue class does not implement 'peek'"): + q.peek() + self.assertEqual(len(q), 3) + self.assertEqual(q.pop().url, req1.url) + self.assertEqual(len(q), 2) + self.assertEqual(q.pop().url, req2.url) + self.assertEqual(len(q), 1) + self.assertEqual(q.pop().url, req3.url) + self.assertEqual(len(q), 0) + self.assertIsNone(q.pop()) + q.close() + + +class LifoQueueMixin(RequestQueueTestMixin): + def test_lifo_with_peek(self): + if not hasattr(queuelib.queue.FifoMemoryQueue, "peek"): + raise unittest.SkipTest("The queuelib queues do not define peek") + q = self.queue() + self.assertEqual(len(q), 0) + self.assertIsNone(q.peek()) + self.assertIsNone(q.pop()) + req1 = Request("http://www.example.com/1") + req2 = Request("http://www.example.com/2") + req3 = Request("http://www.example.com/3") + q.push(req1) + q.push(req2) + q.push(req3) + self.assertEqual(len(q), 3) + self.assertEqual(q.peek().url, req3.url) + self.assertEqual(q.pop().url, req3.url) + self.assertEqual(len(q), 2) + self.assertEqual(q.peek().url, req2.url) + self.assertEqual(q.pop().url, req2.url) + self.assertEqual(len(q), 1) + self.assertEqual(q.peek().url, req1.url) + self.assertEqual(q.pop().url, req1.url) + self.assertEqual(len(q), 0) + self.assertIsNone(q.peek()) + self.assertIsNone(q.pop()) + q.close() + + def test_lifo_without_peek(self): + if hasattr(queuelib.queue.FifoMemoryQueue, "peek"): + raise unittest.SkipTest("The queuelib queues do not define peek") + q = self.queue() + self.assertEqual(len(q), 0) + self.assertIsNone(q.pop()) + req1 = Request("http://www.example.com/1") + req2 = Request("http://www.example.com/2") + req3 = Request("http://www.example.com/3") + q.push(req1) + q.push(req2) + q.push(req3) + with self.assertRaises(NotImplementedError, msg="The underlying queue class does not implement 'peek'"): + q.peek() + self.assertEqual(len(q), 3) + self.assertEqual(q.pop().url, req3.url) + self.assertEqual(len(q), 2) + self.assertEqual(q.pop().url, req2.url) + self.assertEqual(len(q), 1) + self.assertEqual(q.pop().url, req1.url) + self.assertEqual(len(q), 0) + self.assertIsNone(q.pop()) + q.close() + + +class PickleFifoDiskQueueRequestTest(FifoQueueMixin, BaseQueueTestCase): + def queue(self): + return PickleFifoDiskQueue.from_crawler(crawler=self.crawler, key="pickle/fifo") + + +class PickleLifoDiskQueueRequestTest(LifoQueueMixin, BaseQueueTestCase): + def queue(self): + return PickleLifoDiskQueue.from_crawler(crawler=self.crawler, key="pickle/lifo") + + +class MarshalFifoDiskQueueRequestTest(FifoQueueMixin, BaseQueueTestCase): + def queue(self): + return MarshalFifoDiskQueue.from_crawler(crawler=self.crawler, key="marshal/fifo") + + +class MarshalLifoDiskQueueRequestTest(LifoQueueMixin, BaseQueueTestCase): + def queue(self): + return MarshalLifoDiskQueue.from_crawler(crawler=self.crawler, key="marshal/lifo") + + +class FifoMemoryQueueRequestTest(FifoQueueMixin, BaseQueueTestCase): + def queue(self): + return FifoMemoryQueue.from_crawler(crawler=self.crawler) + + +class LifoMemoryQueueRequestTest(LifoQueueMixin, BaseQueueTestCase): + def queue(self): + return LifoMemoryQueue.from_crawler(crawler=self.crawler) From ddea6b7bfa38bf5402d78350ab61e2c827ca49b5 Mon Sep 17 00:00:00 2001 From: Eugenio Lacuesta <1731933+elacuesta@users.noreply.github.com> Date: Mon, 26 Apr 2021 16:16:14 -0300 Subject: [PATCH 08/16] Scheduler: minimal interface, API docs (#3559) --- docs/index.rst | 4 + docs/topics/architecture.rst | 5 +- docs/topics/scheduler.rst | 34 ++++ docs/topics/settings.rst | 3 +- scrapy/core/engine.py | 21 ++- scrapy/core/scheduler.py | 298 +++++++++++++++++++++++++++-------- scrapy/utils/job.py | 5 +- tests/test_scheduler_base.py | 159 +++++++++++++++++++ 8 files changed, 458 insertions(+), 71 deletions(-) create mode 100644 docs/topics/scheduler.rst create mode 100644 tests/test_scheduler_base.py diff --git a/docs/index.rst b/docs/index.rst index da264fb34..433798aa8 100644 --- a/docs/index.rst +++ b/docs/index.rst @@ -227,6 +227,7 @@ Extending Scrapy topics/extensions topics/api topics/signals + topics/scheduler topics/exporters @@ -248,6 +249,9 @@ Extending Scrapy :doc:`topics/signals` See all available signals and how to work with them. +:doc:`topics/scheduler` + Understand the scheduler component. + :doc:`topics/exporters` Quickly export your scraped items to a file (XML, CSV, etc). diff --git a/docs/topics/architecture.rst b/docs/topics/architecture.rst index 074c59241..71d027c86 100644 --- a/docs/topics/architecture.rst +++ b/docs/topics/architecture.rst @@ -87,8 +87,9 @@ of the system, and triggering events when certain actions occur. See the Scheduler --------- -The Scheduler receives requests from the engine and enqueues them for feeding -them later (also to the engine) when the engine requests them. +The :ref:`scheduler ` receives requests from the engine and +enqueues them for feeding them later (also to the engine) when the engine +requests them. .. _component-downloader: diff --git a/docs/topics/scheduler.rst b/docs/topics/scheduler.rst new file mode 100644 index 000000000..57c24b76a --- /dev/null +++ b/docs/topics/scheduler.rst @@ -0,0 +1,34 @@ +.. _topics-scheduler: + +========= +Scheduler +========= + +.. module:: scrapy.core.scheduler + +The scheduler component receives requests from the :ref:`engine ` +and stores them into persistent and/or non-persistent data structures. +It also gets those requests and feeds them back to the engine when it +asks for a next request to be downloaded. + + +Overriding the default scheduler +================================ + +You can use your own custom scheduler class by supplying its full +Python path in the :setting:`SCHEDULER` setting. + + +Minimal scheduler interface +=========================== + +.. autoclass:: BaseScheduler + :members: + + +Default Scrapy scheduler +======================== + +.. autoclass:: Scheduler + :members: + :special-members: __len__ diff --git a/docs/topics/settings.rst b/docs/topics/settings.rst index 1d5babcec..e4fb2baf7 100644 --- a/docs/topics/settings.rst +++ b/docs/topics/settings.rst @@ -1280,7 +1280,8 @@ SCHEDULER Default: ``'scrapy.core.scheduler.Scheduler'`` -The scheduler to use for crawling. +The scheduler class to be used for crawling. +See the :ref:`topics-scheduler` topic for details. .. setting:: SCHEDULER_DEBUG diff --git a/scrapy/core/engine.py b/scrapy/core/engine.py index 7a09bafa1..dd3225082 100644 --- a/scrapy/core/engine.py +++ b/scrapy/core/engine.py @@ -17,6 +17,7 @@ from scrapy import signals from scrapy.core.scraper import Scraper from scrapy.exceptions import DontCloseSpider, ScrapyDeprecationWarning from scrapy.http import Response, Request +from scrapy.settings import BaseSettings from scrapy.spiders import Spider from scrapy.utils.log import logformatter_adapter, failure_to_exc_info from scrapy.utils.misc import create_instance, load_object @@ -73,12 +74,22 @@ class ExecutionEngine: self.spider: Optional[Spider] = None self.running = False self.paused = False - self.scheduler_cls = load_object(crawler.settings["SCHEDULER"]) + self.scheduler_cls = self._get_scheduler_class(crawler.settings) downloader_cls = load_object(self.settings['DOWNLOADER']) self.downloader = downloader_cls(crawler) self.scraper = Scraper(crawler) self._spider_closed_callback = spider_closed_callback + def _get_scheduler_class(self, settings: BaseSettings) -> type: + from scrapy.core.scheduler import BaseScheduler + scheduler_cls = load_object(settings["SCHEDULER"]) + if not issubclass(scheduler_cls, BaseScheduler): + raise TypeError( + f"The provided scheduler class ({settings['SCHEDULER']})" + " does not fully implement the scheduler interface" + ) + return scheduler_cls + @inlineCallbacks def start(self) -> Deferred: if self.running: @@ -301,7 +312,8 @@ class ExecutionEngine: start_requests = yield self.scraper.spidermw.process_start_requests(start_requests, spider) self.slot = Slot(start_requests, close_if_idle, nextcall, scheduler) self.spider = spider - yield scheduler.open(spider) + if hasattr(scheduler, "open"): + yield scheduler.open(spider) yield self.scraper.open_spider(spider) self.crawler.stats.open_spider(spider) yield self.signals.send_catch_log_deferred(signals.spider_opened, spider=spider) @@ -345,8 +357,9 @@ class ExecutionEngine: dfd.addBoth(lambda _: self.scraper.close_spider(spider)) dfd.addErrback(log_failure('Scraper close failure')) - dfd.addBoth(lambda _: self.slot.scheduler.close(reason)) - dfd.addErrback(log_failure('Scheduler close failure')) + if hasattr(self.slot.scheduler, "close"): + dfd.addBoth(lambda _: self.slot.scheduler.close(reason)) + dfd.addErrback(log_failure("Scheduler close failure")) dfd.addBoth(lambda _: self.signals.send_catch_log_deferred( signal=signals.spider_closed, spider=spider, reason=reason, diff --git a/scrapy/core/scheduler.py b/scrapy/core/scheduler.py index 9ce823dbc..5ba0fb63b 100644 --- a/scrapy/core/scheduler.py +++ b/scrapy/core/scheduler.py @@ -1,42 +1,179 @@ -import os import json import logging -from os.path import join, exists +import os +from abc import abstractmethod +from os.path import exists, join +from typing import Optional, Type, TypeVar -from scrapy.utils.misc import load_object, create_instance +from twisted.internet.defer import Deferred + +from scrapy.crawler import Crawler +from scrapy.http.request import Request +from scrapy.spiders import Spider from scrapy.utils.job import job_dir +from scrapy.utils.misc import create_instance, load_object logger = logging.getLogger(__name__) -class Scheduler: +class BaseSchedulerMeta(type): """ - Scrapy Scheduler. It allows to enqueue requests and then get - a next request to download. Scheduler is also handling duplication - filtering, via dupefilter. - - Prioritization and queueing is not performed by the Scheduler. - User sets ``priority`` field for each Request, and a PriorityQueue - (defined by :setting:`SCHEDULER_PRIORITY_QUEUE`) uses these priorities - to dequeue requests in a desired order. - - Scheduler uses two PriorityQueue instances, configured to work in-memory - and on-disk (optional). When on-disk queue is present, it is used by - default, and an in-memory queue is used as a fallback for cases where - a disk queue can't handle a request (can't serialize it). - - :setting:`SCHEDULER_MEMORY_QUEUE` and - :setting:`SCHEDULER_DISK_QUEUE` allow to specify lower-level queue classes - which PriorityQueue instances would be instantiated with, to keep requests - on disk and in memory respectively. - - Overall, Scheduler is an object which holds several PriorityQueue instances - (in-memory and on-disk) and implements fallback logic for them. - Also, it handles dupefilters. + Metaclass to check scheduler classes against the necessary interface """ - def __init__(self, dupefilter, jobdir=None, dqclass=None, mqclass=None, - logunser=False, stats=None, pqclass=None, crawler=None): + def __instancecheck__(cls, instance): + return cls.__subclasscheck__(type(instance)) + + def __subclasscheck__(cls, subclass): + return ( + hasattr(subclass, "has_pending_requests") and callable(subclass.has_pending_requests) + and hasattr(subclass, "enqueue_request") and callable(subclass.enqueue_request) + and hasattr(subclass, "next_request") and callable(subclass.next_request) + ) + + +class BaseScheduler(metaclass=BaseSchedulerMeta): + """ + The scheduler component is responsible for storing requests received from + the engine, and feeding them back upon request (also to the engine). + + The original sources of said requests are: + + * Spider: ``start_requests`` method, requests created for URLs in the ``start_urls`` attribute, request callbacks + * Spider middleware: ``process_spider_output`` and ``process_spider_exception`` methods + * Downloader middleware: ``process_request``, ``process_response`` and ``process_exception`` methods + + The order in which the scheduler returns its stored requests (via the ``next_request`` method) + plays a great part in determining the order in which those requests are downloaded. + + The methods defined in this class constitute the minimal interface that the Scrapy engine will interact with. + """ + + @classmethod + def from_crawler(cls, crawler: Crawler): + """ + Factory method which receives the current :class:`~scrapy.crawler.Crawler` object as argument. + """ + return cls() + + def open(self, spider: Spider) -> Optional[Deferred]: + """ + Called when the spider is opened by the engine. It receives the spider + instance as argument and it's useful to execute initialization code. + + :param spider: the spider object for the current crawl + :type spider: :class:`~scrapy.spiders.Spider` + """ + pass + + def close(self, reason: str) -> Optional[Deferred]: + """ + Called when the spider is closed by the engine. It receives the reason why the crawl + finished as argument and it's useful to execute cleaning code. + + :param reason: a string which describes the reason why the spider was closed + :type reason: :class:`str` + """ + pass + + @abstractmethod + def has_pending_requests(self) -> bool: + """ + ``True`` if the scheduler has enqueued requests, ``False`` otherwise + """ + raise NotImplementedError() + + @abstractmethod + def enqueue_request(self, request: Request) -> bool: + """ + Process a request received by the engine. + + Return ``True`` if the request is stored correctly, ``False`` otherwise. + + If ``False``, the engine will fire a ``request_dropped`` signal, and + will not make further attempts to schedule the request at a later time. + For reference, the default Scrapy scheduler returns ``False`` when the + request is rejected by the dupefilter. + """ + raise NotImplementedError() + + @abstractmethod + def next_request(self) -> Optional[Request]: + """ + Return the next :class:`~scrapy.http.Request` to be processed, or ``None`` + to indicate that there are no requests to be considered ready at the moment. + + Returning ``None`` implies that no request from the scheduler will be sent + to the downloader in the current reactor cycle. The engine will continue + calling ``next_request`` until ``has_pending_requests`` is ``False``. + """ + raise NotImplementedError() + + +SchedulerTV = TypeVar("SchedulerTV", bound="Scheduler") + + +class Scheduler(BaseScheduler): + """ + Default Scrapy scheduler. This implementation also handles duplication + filtering via the :setting:`dupefilter `. + + This scheduler stores requests into several priority queues (defined by the + :setting:`SCHEDULER_PRIORITY_QUEUE` setting). In turn, said priority queues + are backed by either memory or disk based queues (respectively defined by the + :setting:`SCHEDULER_MEMORY_QUEUE` and :setting:`SCHEDULER_DISK_QUEUE` settings). + + Request prioritization is almost entirely delegated to the priority queue. The only + prioritization performed by this scheduler is using the disk-based queue if present + (i.e. if the :setting:`JOBDIR` setting is defined) and falling back to the memory-based + queue if a serialization error occurs. If the disk queue is not present, the memory one + is used directly. + + :param dupefilter: An object responsible for checking and filtering duplicate requests. + The value for the :setting:`DUPEFILTER_CLASS` setting is used by default. + :type dupefilter: :class:`scrapy.dupefilters.BaseDupeFilter` instance or similar: + any class that implements the `BaseDupeFilter` interface + + :param jobdir: The path of a directory to be used for persisting the crawl's state. + The value for the :setting:`JOBDIR` setting is used by default. + See :ref:`topics-jobs`. + :type jobdir: :class:`str` or ``None`` + + :param dqclass: A class to be used as persistent request queue. + The value for the :setting:`SCHEDULER_DISK_QUEUE` setting is used by default. + :type dqclass: class + + :param mqclass: A class to be used as non-persistent request queue. + The value for the :setting:`SCHEDULER_MEMORY_QUEUE` setting is used by default. + :type mqclass: class + + :param logunser: A boolean that indicates whether or not unserializable requests should be logged. + The value for the :setting:`SCHEDULER_DEBUG` setting is used by default. + :type logunser: bool + + :param stats: A stats collector object to record stats about the request scheduling process. + The value for the :setting:`STATS_CLASS` setting is used by default. + :type stats: :class:`scrapy.statscollectors.StatsCollector` instance or similar: + any class that implements the `StatsCollector` interface + + :param pqclass: A class to be used as priority queue for requests. + The value for the :setting:`SCHEDULER_PRIORITY_QUEUE` setting is used by default. + :type pqclass: class + + :param crawler: The crawler object corresponding to the current crawl. + :type crawler: :class:`scrapy.crawler.Crawler` + """ + def __init__( + self, + dupefilter, + jobdir: Optional[str] = None, + dqclass=None, + mqclass=None, + logunser: bool = False, + stats=None, + pqclass=None, + crawler: Optional[Crawler] = None, + ): self.df = dupefilter self.dqdir = self._dqdir(jobdir) self.pqclass = pqclass @@ -47,34 +184,57 @@ class Scheduler: self.crawler = crawler @classmethod - def from_crawler(cls, crawler): - settings = crawler.settings - dupefilter_cls = load_object(settings['DUPEFILTER_CLASS']) - dupefilter = create_instance(dupefilter_cls, settings, crawler) - pqclass = load_object(settings['SCHEDULER_PRIORITY_QUEUE']) - dqclass = load_object(settings['SCHEDULER_DISK_QUEUE']) - mqclass = load_object(settings['SCHEDULER_MEMORY_QUEUE']) - logunser = settings.getbool('SCHEDULER_DEBUG') - return cls(dupefilter, jobdir=job_dir(settings), logunser=logunser, - stats=crawler.stats, pqclass=pqclass, dqclass=dqclass, - mqclass=mqclass, crawler=crawler) + def from_crawler(cls: Type[SchedulerTV], crawler) -> SchedulerTV: + """ + Factory method, initializes the scheduler with arguments taken from the crawl settings + """ + dupefilter_cls = load_object(crawler.settings['DUPEFILTER_CLASS']) + return cls( + dupefilter=create_instance(dupefilter_cls, crawler.settings, crawler), + jobdir=job_dir(crawler.settings), + dqclass=load_object(crawler.settings['SCHEDULER_DISK_QUEUE']), + mqclass=load_object(crawler.settings['SCHEDULER_MEMORY_QUEUE']), + logunser=crawler.settings.getbool('SCHEDULER_DEBUG'), + stats=crawler.stats, + pqclass=load_object(crawler.settings['SCHEDULER_PRIORITY_QUEUE']), + crawler=crawler, + ) - def has_pending_requests(self): + def has_pending_requests(self) -> bool: return len(self) > 0 - def open(self, spider): + def open(self, spider: Spider) -> Optional[Deferred]: + """ + (1) initialize the memory queue + (2) initialize the disk queue if the ``jobdir`` attribute is a valid directory + (3) return the result of the dupefilter's ``open`` method + """ self.spider = spider self.mqs = self._mq() self.dqs = self._dq() if self.dqdir else None return self.df.open() - def close(self, reason): - if self.dqs: + def close(self, reason: str) -> Optional[Deferred]: + """ + (1) dump pending requests to disk if there is a disk queue + (2) return the result of the dupefilter's ``close`` method + """ + if self.dqs is not None: state = self.dqs.close() + assert isinstance(self.dqdir, str) self._write_dqs_state(self.dqdir, state) return self.df.close(reason) - def enqueue_request(self, request): + def enqueue_request(self, request: Request) -> bool: + """ + Unless the received request is filtered out by the Dupefilter, attempt to push + it into the disk queue, falling back to pushing it into the memory queue. + + Increment the appropriate stats, such as: ``scheduler/enqueued``, + ``scheduler/enqueued/disk``, ``scheduler/enqueued/memory``. + + Return ``True`` if the request was stored successfully, ``False`` otherwise. + """ if not request.dont_filter and self.df.request_seen(request): self.df.log(request, self.spider) return False @@ -87,24 +247,35 @@ class Scheduler: self.stats.inc_value('scheduler/enqueued', spider=self.spider) return True - def next_request(self): + def next_request(self) -> Optional[Request]: + """ + Return a :class:`~scrapy.http.Request` object from the memory queue, + falling back to the disk queue if the memory queue is empty. + Return ``None`` if there are no more enqueued requests. + + Increment the appropriate stats, such as: ``scheduler/dequeued``, + ``scheduler/dequeued/disk``, ``scheduler/dequeued/memory``. + """ request = self.mqs.pop() - if request: + if request is not None: self.stats.inc_value('scheduler/dequeued/memory', spider=self.spider) else: request = self._dqpop() - if request: + if request is not None: self.stats.inc_value('scheduler/dequeued/disk', spider=self.spider) - if request: + if request is not None: self.stats.inc_value('scheduler/dequeued', spider=self.spider) return request - def __len__(self): - return len(self.dqs) + len(self.mqs) if self.dqs else len(self.mqs) + def __len__(self) -> int: + """ + Return the total amount of enqueued requests + """ + return len(self.dqs) + len(self.mqs) if self.dqs is not None else len(self.mqs) - def _dqpush(self, request): + def _dqpush(self, request: Request) -> bool: if self.dqs is None: - return + return False try: self.dqs.push(request) except ValueError as e: # non serializable request @@ -115,18 +286,18 @@ class Scheduler: logger.warning(msg, {'request': request, 'reason': e}, exc_info=True, extra={'spider': self.spider}) self.logunser = False - self.stats.inc_value('scheduler/unserializable', - spider=self.spider) - return + self.stats.inc_value('scheduler/unserializable', spider=self.spider) + return False else: return True - def _mqpush(self, request): + def _mqpush(self, request: Request) -> None: self.mqs.push(request) - def _dqpop(self): - if self.dqs: + def _dqpop(self) -> Optional[Request]: + if self.dqs is not None: return self.dqs.pop() + return None def _mq(self): """ Create a new priority queue instance, with in-memory storage """ @@ -150,21 +321,22 @@ class Scheduler: {'queuesize': len(q)}, extra={'spider': self.spider}) return q - def _dqdir(self, jobdir): + def _dqdir(self, jobdir: Optional[str]) -> Optional[str]: """ Return a folder name to keep disk queue state at """ - if jobdir: + if jobdir is not None: dqdir = join(jobdir, 'requests.queue') if not exists(dqdir): os.makedirs(dqdir) return dqdir + return None - def _read_dqs_state(self, dqdir): + def _read_dqs_state(self, dqdir: str) -> list: path = join(dqdir, 'active.json') if not exists(path): - return () + return [] with open(path) as f: return json.load(f) - def _write_dqs_state(self, dqdir, state): + def _write_dqs_state(self, dqdir: str, state: list) -> None: with open(join(dqdir, 'active.json'), 'w') as f: json.dump(state, f) diff --git a/scrapy/utils/job.py b/scrapy/utils/job.py index 4f1e601fc..c92ef36f5 100644 --- a/scrapy/utils/job.py +++ b/scrapy/utils/job.py @@ -1,7 +1,10 @@ import os +from typing import Optional + +from scrapy.settings import BaseSettings -def job_dir(settings): +def job_dir(settings: BaseSettings) -> Optional[str]: path = settings['JOBDIR'] if path and not os.path.exists(path): os.makedirs(path) diff --git a/tests/test_scheduler_base.py b/tests/test_scheduler_base.py new file mode 100644 index 000000000..bf90b4320 --- /dev/null +++ b/tests/test_scheduler_base.py @@ -0,0 +1,159 @@ +from typing import Dict, Optional +from unittest import TestCase +from urllib.parse import urljoin, urlparse + +from testfixtures import LogCapture +from twisted.internet import defer +from twisted.trial.unittest import TestCase as TwistedTestCase + +from scrapy.core.scheduler import BaseScheduler +from scrapy.crawler import CrawlerRunner +from scrapy.http import Request +from scrapy.spiders import Spider +from scrapy.utils.request import request_fingerprint + +from tests.mockserver import MockServer + + +PATHS = ["/a", "/b", "/c"] +URLS = [urljoin("https://example.org", p) for p in PATHS] + + +class MinimalScheduler: + def __init__(self) -> None: + self.requests: Dict[str, Request] = {} + + def has_pending_requests(self) -> bool: + return bool(self.requests) + + def enqueue_request(self, request: Request) -> bool: + fp = request_fingerprint(request) + if fp not in self.requests: + self.requests[fp] = request + return True + return False + + def next_request(self) -> Optional[Request]: + if self.has_pending_requests(): + fp, request = self.requests.popitem() + return request + return None + + +class SimpleScheduler(MinimalScheduler): + def open(self, spider: Spider) -> defer.Deferred: + return defer.succeed("open") + + def close(self, reason: str) -> defer.Deferred: + return defer.succeed("close") + + def __len__(self) -> int: + return len(self.requests) + + +class TestSpider(Spider): + name = "test" + + def __init__(self, mockserver, *args, **kwargs): + super().__init__(*args, **kwargs) + self.start_urls = map(mockserver.url, PATHS) + + def parse(self, response): + return {"path": urlparse(response.url).path} + + +class InterfaceCheckMixin: + def test_scheduler_class(self): + self.assertTrue(isinstance(self.scheduler, BaseScheduler)) + self.assertTrue(issubclass(self.scheduler.__class__, BaseScheduler)) + + +class BaseSchedulerTest(TestCase, InterfaceCheckMixin): + def setUp(self): + self.scheduler = BaseScheduler() + + def test_methods(self): + self.assertIsNone(self.scheduler.open(Spider("foo"))) + self.assertIsNone(self.scheduler.close("finished")) + self.assertRaises(NotImplementedError, self.scheduler.has_pending_requests) + self.assertRaises(NotImplementedError, self.scheduler.enqueue_request, Request("https://example.org")) + self.assertRaises(NotImplementedError, self.scheduler.next_request) + + +class MinimalSchedulerTest(TestCase, InterfaceCheckMixin): + def setUp(self): + self.scheduler = MinimalScheduler() + + def test_open_close(self): + with self.assertRaises(AttributeError): + self.scheduler.open(Spider("foo")) + with self.assertRaises(AttributeError): + self.scheduler.close("finished") + + def test_len(self): + with self.assertRaises(AttributeError): + self.scheduler.__len__() + with self.assertRaises(TypeError): + len(self.scheduler) + + def test_enqueue_dequeue(self): + self.assertFalse(self.scheduler.has_pending_requests()) + for url in URLS: + self.assertTrue(self.scheduler.enqueue_request(Request(url))) + self.assertFalse(self.scheduler.enqueue_request(Request(url))) + self.assertTrue(self.scheduler.has_pending_requests) + + dequeued = [] + while self.scheduler.has_pending_requests(): + request = self.scheduler.next_request() + dequeued.append(request.url) + self.assertEqual(set(dequeued), set(URLS)) + self.assertFalse(self.scheduler.has_pending_requests()) + + +class SimpleSchedulerTest(TwistedTestCase, InterfaceCheckMixin): + def setUp(self): + self.scheduler = SimpleScheduler() + + @defer.inlineCallbacks + def test_enqueue_dequeue(self): + open_result = yield self.scheduler.open(Spider("foo")) + self.assertEqual(open_result, "open") + self.assertFalse(self.scheduler.has_pending_requests()) + + for url in URLS: + self.assertTrue(self.scheduler.enqueue_request(Request(url))) + self.assertFalse(self.scheduler.enqueue_request(Request(url))) + + self.assertTrue(self.scheduler.has_pending_requests()) + self.assertEqual(len(self.scheduler), len(URLS)) + + dequeued = [] + while self.scheduler.has_pending_requests(): + request = self.scheduler.next_request() + dequeued.append(request.url) + self.assertEqual(set(dequeued), set(URLS)) + + self.assertFalse(self.scheduler.has_pending_requests()) + self.assertEqual(len(self.scheduler), 0) + + close_result = yield self.scheduler.close("") + self.assertEqual(close_result, "close") + + +class MinimalSchedulerCrawlTest(TwistedTestCase): + scheduler_cls = MinimalScheduler + + @defer.inlineCallbacks + def test_crawl(self): + with MockServer() as mockserver: + settings = {"SCHEDULER": self.scheduler_cls} + with LogCapture() as log: + yield CrawlerRunner(settings).crawl(TestSpider, mockserver) + for path in PATHS: + self.assertIn(f"{{'path': '{path}'}}", str(log)) + self.assertIn(f"'item_scraped_count': {len(PATHS)}", str(log)) + + +class SimpleSchedulerCrawlTest(MinimalSchedulerCrawlTest): + scheduler_cls = SimpleScheduler From 02ae1deaf499e79dc8ececf4d5dcea6daef0ade0 Mon Sep 17 00:00:00 2001 From: Eugenio Lacuesta <1731933+elacuesta@users.noreply.github.com> Date: Tue, 27 Apr 2021 09:41:44 -0300 Subject: [PATCH 09/16] Deprecate unused squeues (#5117) --- scrapy/squeues.py | 47 +++++++++++++++++++++++++++++++++++-------- tests/test_squeues.py | 16 +++++++-------- 2 files changed, 47 insertions(+), 16 deletions(-) diff --git a/scrapy/squeues.py b/scrapy/squeues.py index 44898ba08..16f7bf4b6 100644 --- a/scrapy/squeues.py +++ b/scrapy/squeues.py @@ -8,6 +8,7 @@ import pickle from queuelib import queue +from scrapy.utils.deprecate import create_deprecated_class from scrapy.utils.reqser import request_to_dict, request_from_dict @@ -123,30 +124,60 @@ def _pickle_serialize(obj): raise ValueError(str(e)) from e -PickleFifoDiskQueueNonRequest = _serializable_queue( +_PickleFifoSerializationDiskQueue = _serializable_queue( _with_mkdir(queue.FifoDiskQueue), _pickle_serialize, pickle.loads ) -PickleLifoDiskQueueNonRequest = _serializable_queue( +_PickleLifoSerializationDiskQueue = _serializable_queue( _with_mkdir(queue.LifoDiskQueue), _pickle_serialize, pickle.loads ) -MarshalFifoDiskQueueNonRequest = _serializable_queue( +_MarshalFifoSerializationDiskQueue = _serializable_queue( _with_mkdir(queue.FifoDiskQueue), marshal.dumps, marshal.loads ) -MarshalLifoDiskQueueNonRequest = _serializable_queue( +_MarshalLifoSerializationDiskQueue = _serializable_queue( _with_mkdir(queue.LifoDiskQueue), marshal.dumps, marshal.loads ) -PickleFifoDiskQueue = _scrapy_serialization_queue(PickleFifoDiskQueueNonRequest) -PickleLifoDiskQueue = _scrapy_serialization_queue(PickleLifoDiskQueueNonRequest) -MarshalFifoDiskQueue = _scrapy_serialization_queue(MarshalFifoDiskQueueNonRequest) -MarshalLifoDiskQueue = _scrapy_serialization_queue(MarshalLifoDiskQueueNonRequest) +# public queue classes +PickleFifoDiskQueue = _scrapy_serialization_queue(_PickleFifoSerializationDiskQueue) +PickleLifoDiskQueue = _scrapy_serialization_queue(_PickleLifoSerializationDiskQueue) +MarshalFifoDiskQueue = _scrapy_serialization_queue(_MarshalFifoSerializationDiskQueue) +MarshalLifoDiskQueue = _scrapy_serialization_queue(_MarshalLifoSerializationDiskQueue) FifoMemoryQueue = _scrapy_non_serialization_queue(queue.FifoMemoryQueue) LifoMemoryQueue = _scrapy_non_serialization_queue(queue.LifoMemoryQueue) + + +# deprecated queue classes +_subclass_warn_message = "{cls} inherits from deprecated class {old}" +_instance_warn_message = "{cls} is deprecated" +PickleFifoDiskQueueNonRequest = create_deprecated_class( + name="PickleFifoDiskQueueNonRequest", + new_class=_PickleFifoSerializationDiskQueue, + subclass_warn_message=_subclass_warn_message, + instance_warn_message=_instance_warn_message, +) +PickleLifoDiskQueueNonRequest = create_deprecated_class( + name="PickleLifoDiskQueueNonRequest", + new_class=_PickleLifoSerializationDiskQueue, + subclass_warn_message=_subclass_warn_message, + instance_warn_message=_instance_warn_message, +) +MarshalFifoDiskQueueNonRequest = create_deprecated_class( + name="MarshalFifoDiskQueueNonRequest", + new_class=_MarshalFifoSerializationDiskQueue, + subclass_warn_message=_subclass_warn_message, + instance_warn_message=_instance_warn_message, +) +MarshalLifoDiskQueueNonRequest = create_deprecated_class( + name="MarshalLifoDiskQueueNonRequest", + new_class=_MarshalLifoSerializationDiskQueue, + subclass_warn_message=_subclass_warn_message, + instance_warn_message=_instance_warn_message, +) diff --git a/tests/test_squeues.py b/tests/test_squeues.py index becacce62..acc821b83 100644 --- a/tests/test_squeues.py +++ b/tests/test_squeues.py @@ -3,10 +3,10 @@ import sys from queuelib.tests import test_queue as t from scrapy.squeues import ( - MarshalFifoDiskQueueNonRequest as MarshalFifoDiskQueue, - MarshalLifoDiskQueueNonRequest as MarshalLifoDiskQueue, - PickleFifoDiskQueueNonRequest as PickleFifoDiskQueue, - PickleLifoDiskQueueNonRequest as PickleLifoDiskQueue + _MarshalFifoSerializationDiskQueue, + _MarshalLifoSerializationDiskQueue, + _PickleFifoSerializationDiskQueue, + _PickleLifoSerializationDiskQueue, ) from scrapy.item import Item, Field from scrapy.http import Request @@ -53,7 +53,7 @@ class MarshalFifoDiskQueueTest(t.FifoDiskQueueTest, FifoDiskQueueTestMixin): chunksize = 100000 def queue(self): - return MarshalFifoDiskQueue(self.qpath, chunksize=self.chunksize) + return _MarshalFifoSerializationDiskQueue(self.qpath, chunksize=self.chunksize) class ChunkSize1MarshalFifoDiskQueueTest(MarshalFifoDiskQueueTest): @@ -77,7 +77,7 @@ class PickleFifoDiskQueueTest(t.FifoDiskQueueTest, FifoDiskQueueTestMixin): chunksize = 100000 def queue(self): - return PickleFifoDiskQueue(self.qpath, chunksize=self.chunksize) + return _PickleFifoSerializationDiskQueue(self.qpath, chunksize=self.chunksize) def test_serialize_item(self): q = self.queue() @@ -155,13 +155,13 @@ class LifoDiskQueueTestMixin: class MarshalLifoDiskQueueTest(t.LifoDiskQueueTest, LifoDiskQueueTestMixin): def queue(self): - return MarshalLifoDiskQueue(self.qpath) + return _MarshalLifoSerializationDiskQueue(self.qpath) class PickleLifoDiskQueueTest(t.LifoDiskQueueTest, LifoDiskQueueTestMixin): def queue(self): - return PickleLifoDiskQueue(self.qpath) + return _PickleLifoSerializationDiskQueue(self.qpath) def test_serialize_item(self): q = self.queue() From 4f500342c8ad4674b191e1fab0d1b2ac944d7d3e Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Tom=C3=A1=C5=A1=20Hrn=C4=8Diar?= Date: Wed, 28 Apr 2021 11:57:44 +0200 Subject: [PATCH 10/16] Require setuptools, scrapy/cmdline.py, /setup.py and tests/test_webclient.py import pkg_resources --- setup.py | 1 + 1 file changed, 1 insertion(+) diff --git a/setup.py b/setup.py index 2b60a10af..b1bb64575 100644 --- a/setup.py +++ b/setup.py @@ -32,6 +32,7 @@ install_requires = [ 'protego>=0.1.15', 'itemadapter>=0.1.0', 'h2>=3.0,<4.0', + 'setuptools', ] extras_require = {} cpython_dependencies = [ From 34b216289c31c27ad6256ff95382505a9d84adb3 Mon Sep 17 00:00:00 2001 From: Renne Rocha Date: Thu, 6 May 2021 11:34:05 -0300 Subject: [PATCH 11/16] Update link for reasoning value of URLLENGTH_LIMIT (#5134) --- docs/topics/settings.rst | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/docs/topics/settings.rst b/docs/topics/settings.rst index e4fb2baf7..2506497e2 100644 --- a/docs/topics/settings.rst +++ b/docs/topics/settings.rst @@ -1619,7 +1619,7 @@ Default: ``2083`` Scope: ``spidermiddlewares.urllength`` The maximum URL length to allow for crawled URLs. For more information about -the default value for this setting see: https://boutell.com/newfaq/misc/urllength.html +the default value for this setting see: https://support.microsoft.com/en-us/topic/maximum-url-length-is-2-083-characters-in-internet-explorer-174e7c8a-6666-f4e0-6fd6-908b53c12246 .. setting:: USER_AGENT @@ -1642,7 +1642,6 @@ case to see how to enable and use them. .. settingslist:: - .. _Amazon web services: https://aws.amazon.com/ .. _breadth-first order: https://en.wikipedia.org/wiki/Breadth-first_search .. _depth-first order: https://en.wikipedia.org/wiki/Depth-first_search From cec36a9284641bb7b69a2081ad92cd7d4ee25934 Mon Sep 17 00:00:00 2001 From: Eugenio Lacuesta <1731933+elacuesta@users.noreply.github.com> Date: Mon, 10 May 2021 13:00:08 -0300 Subject: [PATCH 12/16] Refactor request to/from dict (#5130) --- docs/topics/request-response.rst | 17 ++- scrapy/http/request/__init__.py | 70 ++++++++++-- scrapy/http/request/json_request.py | 8 ++ scrapy/squeues.py | 8 +- scrapy/utils/reqser.py | 103 +++--------------- scrapy/utils/request.py | 27 ++++- ...t_utils_reqser.py => test_request_dict.py} | 76 +++++++++---- 7 files changed, 183 insertions(+), 126 deletions(-) rename tests/{test_utils_reqser.py => test_request_dict.py} (68%) diff --git a/docs/topics/request-response.rst b/docs/topics/request-response.rst index 500781c05..73b5a858f 100644 --- a/docs/topics/request-response.rst +++ b/docs/topics/request-response.rst @@ -26,10 +26,6 @@ Request objects .. autoclass:: Request - A :class:`Request` object represents an HTTP request, which is usually - generated in the Spider and executed by the Downloader, and thus generating - a :class:`Response`. - :param url: the URL of this request If the URL is invalid, a :exc:`ValueError` exception is raised. @@ -205,6 +201,8 @@ Request objects ``failure.request.cb_kwargs`` in the request's errback. For more information, see :ref:`errback-cb_kwargs`. + .. autoattribute:: Request.attributes + .. method:: Request.copy() Return a new Request which is a copy of this Request. See also: @@ -220,6 +218,15 @@ Request objects .. automethod:: from_curl + .. automethod:: to_dict + + +Other functions related to requests +----------------------------------- + +.. autofunction:: scrapy.utils.request.request_from_dict + + .. _topics-request-response-ref-request-callback-arguments: Passing additional data to callback functions @@ -642,6 +649,8 @@ dealing with JSON requests. data into JSON format. :type dumps_kwargs: dict + .. autoattribute:: JsonRequest.attributes + JsonRequest usage example ------------------------- diff --git a/scrapy/http/request/__init__.py b/scrapy/http/request/__init__.py index 498f1b052..ad884feac 100644 --- a/scrapy/http/request/__init__.py +++ b/scrapy/http/request/__init__.py @@ -4,17 +4,37 @@ requests in Scrapy. See documentation in docs/topics/request-response.rst """ +import inspect +from typing import Optional, Tuple + from w3lib.url import safe_url_string +import scrapy +from scrapy.http.common import obsolete_setter from scrapy.http.headers import Headers +from scrapy.utils.curl import curl_to_request_kwargs from scrapy.utils.python import to_bytes from scrapy.utils.trackref import object_ref from scrapy.utils.url import escape_ajax -from scrapy.http.common import obsolete_setter -from scrapy.utils.curl import curl_to_request_kwargs class Request(object_ref): + """Represents an HTTP request, which is usually generated in a Spider and + executed by the Downloader, thus generating a :class:`Response`. + """ + + attributes: Tuple[str, ...] = ( + "url", "callback", "method", "headers", "body", + "cookies", "meta", "encoding", "priority", + "dont_filter", "errback", "flags", "cb_kwargs", + ) + """A tuple of :class:`str` objects containing the name of all public + attributes of the class that are also keyword parameters of the + ``__init__`` method. + + Currently used by :meth:`Request.replace`, :meth:`Request.to_dict` and + :func:`~scrapy.utils.request.request_from_dict`. + """ def __init__(self, url, callback=None, method='GET', headers=None, body=None, cookies=None, meta=None, encoding='utf-8', priority=0, @@ -99,11 +119,8 @@ class Request(object_ref): return self.replace() def replace(self, *args, **kwargs): - """Create a new Request with the same attributes except for those - given new values. - """ - for x in ['url', 'method', 'headers', 'body', 'cookies', 'meta', 'flags', - 'encoding', 'priority', 'dont_filter', 'callback', 'errback', 'cb_kwargs']: + """Create a new Request with the same attributes except for those given new values""" + for x in self.attributes: kwargs.setdefault(x, getattr(self, x)) cls = kwargs.pop('cls', self.__class__) return cls(*args, **kwargs) @@ -136,8 +153,43 @@ class Request(object_ref): To translate a cURL command into a Scrapy request, you may use `curl2scrapy `_. - - """ + """ request_kwargs = curl_to_request_kwargs(curl_command, ignore_unknown_options) request_kwargs.update(kwargs) return cls(**request_kwargs) + + def to_dict(self, *, spider: Optional["scrapy.Spider"] = None) -> dict: + """Return a dictionary containing the Request's data. + + Use :func:`~scrapy.utils.request.request_from_dict` to convert back into a :class:`~scrapy.Request` object. + + If a spider is given, this method will try to find out the name of the spider methods used as callback + and errback and include them in the output dict, raising an exception if they cannot be found. + """ + d = { + "url": self.url, # urls are safe (safe_string_url) + "callback": _find_method(spider, self.callback) if callable(self.callback) else self.callback, + "errback": _find_method(spider, self.errback) if callable(self.errback) else self.errback, + "headers": dict(self.headers), + } + for attr in self.attributes: + d.setdefault(attr, getattr(self, attr)) + if type(self) is not Request: + d["_class"] = self.__module__ + '.' + self.__class__.__name__ + return d + + +def _find_method(obj, func): + """Helper function for Request.to_dict""" + # Only instance methods contain ``__func__`` + if obj and hasattr(func, '__func__'): + members = inspect.getmembers(obj, predicate=inspect.ismethod) + for name, obj_func in members: + # We need to use __func__ to access the original function object because instance + # method objects are generated each time attribute is retrieved from instance. + # + # Reference: The standard type hierarchy + # https://docs.python.org/3/reference/datamodel.html + if obj_func.__func__ is func.__func__: + return name + raise ValueError(f"Function {func} is not an instance method in: {obj}") diff --git a/scrapy/http/request/json_request.py b/scrapy/http/request/json_request.py index eae3f9f6b..04e80d897 100644 --- a/scrapy/http/request/json_request.py +++ b/scrapy/http/request/json_request.py @@ -8,12 +8,16 @@ See documentation in docs/topics/request-response.rst import copy import json import warnings +from typing import Tuple from scrapy.http.request import Request from scrapy.utils.deprecate import create_deprecated_class class JsonRequest(Request): + + attributes: Tuple[str, ...] = Request.attributes + ("dumps_kwargs",) + def __init__(self, *args, **kwargs): dumps_kwargs = copy.deepcopy(kwargs.pop('dumps_kwargs', {})) dumps_kwargs.setdefault('sort_keys', True) @@ -36,6 +40,10 @@ class JsonRequest(Request): self.headers.setdefault('Content-Type', 'application/json') self.headers.setdefault('Accept', 'application/json, text/javascript, */*; q=0.01') + @property + def dumps_kwargs(self): + return self._dumps_kwargs + def replace(self, *args, **kwargs): body_passed = kwargs.get('body', None) is not None data = kwargs.pop('data', None) diff --git a/scrapy/squeues.py b/scrapy/squeues.py index 16f7bf4b6..dff9b1350 100644 --- a/scrapy/squeues.py +++ b/scrapy/squeues.py @@ -9,7 +9,7 @@ import pickle from queuelib import queue from scrapy.utils.deprecate import create_deprecated_class -from scrapy.utils.reqser import request_to_dict, request_from_dict +from scrapy.utils.request import request_from_dict def _with_mkdir(queue_class): @@ -68,14 +68,14 @@ def _scrapy_serialization_queue(queue_class): return cls(crawler, key) def push(self, request): - request = request_to_dict(request, self.spider) + request = request.to_dict(spider=self.spider) return super().push(request) def pop(self): request = super().pop() if not request: return None - return request_from_dict(request, self.spider) + return request_from_dict(request, spider=self.spider) def peek(self): """Returns the next object to be returned by :meth:`pop`, @@ -87,7 +87,7 @@ def _scrapy_serialization_queue(queue_class): request = super().peek() if not request: return None - return request_from_dict(request, self.spider) + return request_from_dict(request, spider=self.spider) return ScrapyRequestQueue diff --git a/scrapy/utils/reqser.py b/scrapy/utils/reqser.py index d38b1bc4d..c254b9f82 100644 --- a/scrapy/utils/reqser.py +++ b/scrapy/utils/reqser.py @@ -1,95 +1,22 @@ -""" -Helper functions for serializing (and deserializing) requests. -""" -import inspect +import warnings +from typing import Optional -from scrapy.http import Request -from scrapy.utils.python import to_unicode -from scrapy.utils.misc import load_object +import scrapy +from scrapy.exceptions import ScrapyDeprecationWarning +from scrapy.utils.request import request_from_dict as _from_dict -def request_to_dict(request, spider=None): - """Convert Request object to a dict. - - If a spider is given, it will try to find out the name of the spider method - used in the callback and store that as the callback. - """ - cb = request.callback - if callable(cb): - cb = _find_method(spider, cb) - eb = request.errback - if callable(eb): - eb = _find_method(spider, eb) - d = { - 'url': to_unicode(request.url), # urls should be safe (safe_string_url) - 'callback': cb, - 'errback': eb, - 'method': request.method, - 'headers': dict(request.headers), - 'body': request.body, - 'cookies': request.cookies, - 'meta': request.meta, - '_encoding': request._encoding, - 'priority': request.priority, - 'dont_filter': request.dont_filter, - 'flags': request.flags, - 'cb_kwargs': request.cb_kwargs, - } - if type(request) is not Request: - d['_class'] = request.__module__ + '.' + request.__class__.__name__ - return d +warnings.warn( + ("Module scrapy.utils.reqser is deprecated, please use request.to_dict method" + " and/or scrapy.utils.request.request_from_dict instead"), + category=ScrapyDeprecationWarning, + stacklevel=2, +) -def request_from_dict(d, spider=None): - """Create Request object from a dict. - - If a spider is given, it will try to resolve the callbacks looking at the - spider for methods with the same name. - """ - cb = d['callback'] - if cb and spider: - cb = _get_method(spider, cb) - eb = d['errback'] - if eb and spider: - eb = _get_method(spider, eb) - request_cls = load_object(d['_class']) if '_class' in d else Request - return request_cls( - url=to_unicode(d['url']), - callback=cb, - errback=eb, - method=d['method'], - headers=d['headers'], - body=d['body'], - cookies=d['cookies'], - meta=d['meta'], - encoding=d['_encoding'], - priority=d['priority'], - dont_filter=d['dont_filter'], - flags=d.get('flags'), - cb_kwargs=d.get('cb_kwargs'), - ) +def request_to_dict(request: "scrapy.Request", spider: Optional["scrapy.Spider"] = None) -> dict: + return request.to_dict(spider=spider) -def _find_method(obj, func): - # Only instance methods contain ``__func__`` - if obj and hasattr(func, '__func__'): - members = inspect.getmembers(obj, predicate=inspect.ismethod) - for name, obj_func in members: - # We need to use __func__ to access the original - # function object because instance method objects - # are generated each time attribute is retrieved from - # instance. - # - # Reference: The standard type hierarchy - # https://docs.python.org/3/reference/datamodel.html - if obj_func.__func__ is func.__func__: - return name - raise ValueError(f"Function {func} is not an instance method in: {obj}") - - -def _get_method(obj, name): - name = str(name) - try: - return getattr(obj, name) - except AttributeError: - raise ValueError(f"Method {name!r} not found in: {obj}") +def request_from_dict(d: dict, spider: Optional["scrapy.Spider"] = None) -> "scrapy.Request": + return _from_dict(d, spider=spider) diff --git a/scrapy/utils/request.py b/scrapy/utils/request.py index 541368423..57dcc5f2c 100644 --- a/scrapy/utils/request.py +++ b/scrapy/utils/request.py @@ -11,8 +11,9 @@ from weakref import WeakKeyDictionary from w3lib.http import basic_auth_header from w3lib.url import canonicalize_url -from scrapy.http import Request +from scrapy import Request, Spider from scrapy.utils.httpobj import urlparse_cached +from scrapy.utils.misc import load_object from scrapy.utils.python import to_bytes, to_unicode @@ -106,3 +107,27 @@ def referer_str(request: Request) -> Optional[str]: if referrer is None: return referrer return to_unicode(referrer, errors='replace') + + +def request_from_dict(d: dict, *, spider: Optional[Spider] = None) -> Request: + """Create a :class:`~scrapy.Request` object from a dict. + + If a spider is given, it will try to resolve the callbacks looking at the + spider for methods with the same name. + """ + request_cls = load_object(d["_class"]) if "_class" in d else Request + kwargs = {key: value for key, value in d.items() if key in request_cls.attributes} + if d.get("callback") and spider: + kwargs["callback"] = _get_method(spider, d["callback"]) + if d.get("errback") and spider: + kwargs["errback"] = _get_method(spider, d["errback"]) + return request_cls(**kwargs) + + +def _get_method(obj, name): + """Helper function for request_from_dict""" + name = str(name) + try: + return getattr(obj, name) + except AttributeError: + raise ValueError(f"Method {name!r} not found in: {obj}") diff --git a/tests/test_utils_reqser.py b/tests/test_request_dict.py similarity index 68% rename from tests/test_utils_reqser.py rename to tests/test_request_dict.py index ee68cf6b1..5bdcb975b 100644 --- a/tests/test_utils_reqser.py +++ b/tests/test_request_dict.py @@ -1,8 +1,16 @@ +import sys import unittest +import warnings +from contextlib import suppress -from scrapy.http import Request, FormRequest -from scrapy.spiders import Spider -from scrapy.utils.reqser import request_to_dict, request_from_dict +from scrapy import Spider, Request +from scrapy.exceptions import ScrapyDeprecationWarning +from scrapy.http import FormRequest, JsonRequest +from scrapy.utils.request import request_from_dict + + +class CustomRequest(Request): + pass class RequestSerializationTest(unittest.TestCase): @@ -27,7 +35,8 @@ class RequestSerializationTest(unittest.TestCase): priority=20, meta={'a': 'b'}, cb_kwargs={'k': 'v'}, - flags=['testFlag']) + flags=['testFlag'], + ) self._assert_serializes_ok(r, spider=self.spider) def test_latin1_body(self): @@ -39,7 +48,7 @@ class RequestSerializationTest(unittest.TestCase): self._assert_serializes_ok(r) def _assert_serializes_ok(self, request, spider=None): - d = request_to_dict(request, spider=spider) + d = request.to_dict(spider=spider) request2 = request_from_dict(d, spider=spider) self._assert_same_request(request, request2) @@ -54,16 +63,21 @@ class RequestSerializationTest(unittest.TestCase): self.assertEqual(r1.cookies, r2.cookies) self.assertEqual(r1.meta, r2.meta) self.assertEqual(r1.cb_kwargs, r2.cb_kwargs) + self.assertEqual(r1.encoding, r2.encoding) self.assertEqual(r1._encoding, r2._encoding) self.assertEqual(r1.priority, r2.priority) self.assertEqual(r1.dont_filter, r2.dont_filter) self.assertEqual(r1.flags, r2.flags) + if isinstance(r1, JsonRequest): + self.assertEqual(r1.dumps_kwargs, r2.dumps_kwargs) def test_request_class(self): - r = FormRequest("http://www.example.com") - self._assert_serializes_ok(r, spider=self.spider) - r = CustomRequest("http://www.example.com") - self._assert_serializes_ok(r, spider=self.spider) + r1 = FormRequest("http://www.example.com") + self._assert_serializes_ok(r1, spider=self.spider) + r2 = CustomRequest("http://www.example.com") + self._assert_serializes_ok(r2, spider=self.spider) + r3 = JsonRequest("http://www.example.com", dumps_kwargs={"indent": 4}) + self._assert_serializes_ok(r3, spider=self.spider) def test_callback_serialization(self): r = Request("http://www.example.com", callback=self.spider.parse_item, @@ -75,7 +89,7 @@ class RequestSerializationTest(unittest.TestCase): callback=self.spider.parse_item_reference, errback=self.spider.handle_error_reference) self._assert_serializes_ok(r, spider=self.spider) - request_dict = request_to_dict(r, self.spider) + request_dict = r.to_dict(spider=self.spider) self.assertEqual(request_dict['callback'], 'parse_item_reference') self.assertEqual(request_dict['errback'], 'handle_error_reference') @@ -84,7 +98,7 @@ class RequestSerializationTest(unittest.TestCase): callback=self.spider._TestSpider__parse_item_reference, errback=self.spider._TestSpider__handle_error_reference) self._assert_serializes_ok(r, spider=self.spider) - request_dict = request_to_dict(r, self.spider) + request_dict = r.to_dict(spider=self.spider) self.assertEqual(request_dict['callback'], '_TestSpider__parse_item_reference') self.assertEqual(request_dict['errback'], @@ -110,18 +124,16 @@ class RequestSerializationTest(unittest.TestCase): def test_unserializable_callback1(self): r = Request("http://www.example.com", callback=lambda x: x) - self.assertRaises(ValueError, request_to_dict, r) - self.assertRaises(ValueError, request_to_dict, r, spider=self.spider) + self.assertRaises(ValueError, r.to_dict, spider=self.spider) def test_unserializable_callback2(self): r = Request("http://www.example.com", callback=self.spider.parse_item) - self.assertRaises(ValueError, request_to_dict, r) + self.assertRaises(ValueError, r.to_dict, spider=None) def test_unserializable_callback3(self): """Parser method is removed or replaced dynamically.""" class MySpider(Spider): - name = 'my_spider' def parse(self, response): @@ -130,7 +142,35 @@ class RequestSerializationTest(unittest.TestCase): spider = MySpider() r = Request("http://www.example.com", callback=spider.parse) setattr(spider, 'parse', None) - self.assertRaises(ValueError, request_to_dict, r, spider=spider) + self.assertRaises(ValueError, r.to_dict, spider=spider) + + def test_callback_not_available(self): + """Callback method is not available in the spider passed to from_dict""" + spider = TestSpiderDelegation() + r = Request("http://www.example.com", callback=spider.delegated_callback) + d = r.to_dict(spider=spider) + self.assertRaises(ValueError, request_from_dict, d, spider=Spider("foo")) + + +class DeprecatedMethodsRequestSerializationTest(RequestSerializationTest): + def _assert_serializes_ok(self, request, spider=None): + with warnings.catch_warnings(record=True) as caught: + warnings.simplefilter("always") + with suppress(KeyError): + del sys.modules["scrapy.utils.reqser"] # delete module to reset the deprecation warning + + from scrapy.utils.reqser import request_from_dict as _from_dict, request_to_dict as _to_dict + + request_copy = _from_dict(_to_dict(request, spider), spider) + self._assert_same_request(request, request_copy) + + self.assertEqual(len(caught), 1) + self.assertTrue(issubclass(caught[0].category, ScrapyDeprecationWarning)) + self.assertEqual( + "Module scrapy.utils.reqser is deprecated, please use request.to_dict method" + " and/or scrapy.utils.request.request_from_dict instead", + str(caught[0].message), + ) class TestSpiderMixin: @@ -177,7 +217,3 @@ class TestSpider(Spider, TestSpiderMixin): def __parse_item_private(self, response): pass - - -class CustomRequest(Request): - pass From bd60c3f41fd25e7b5a413cf5c112166b813c2691 Mon Sep 17 00:00:00 2001 From: Shinichi Takayanagi Date: Tue, 11 May 2021 04:58:04 +0900 Subject: [PATCH 13/16] More documentation for setting spider atributes MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit * docs: require sphinx-rtd-theme>=0.5.2 and the latest pip to prevent installing breaking docutils>=0.17 * Update feed-exports.rst * Update feed-exports.rst * Reflects the comments * Remove redundant newline * Update docs/topics/feed-exports.rst Co-authored-by: Adrián Chaves * Apply suggestions from code review Co-authored-by: Adrián Chaves Co-authored-by: Adrián Chaves Co-authored-by: Eugenio Lacuesta --- docs/topics/feed-exports.rst | 3 +++ docs/topics/spiders.rst | 8 ++++++++ 2 files changed, 11 insertions(+) diff --git a/docs/topics/feed-exports.rst b/docs/topics/feed-exports.rst index e772a461c..26c247cdd 100644 --- a/docs/topics/feed-exports.rst +++ b/docs/topics/feed-exports.rst @@ -135,6 +135,9 @@ Here are some examples to illustrate: - ``s3://mybucket/scraping/feeds/%(name)s/%(time)s.json`` +.. note:: :ref:`Spider arguments ` become spider attributes, hence + they can also be used as storage URI parameters. + .. _topics-feed-storage-backends: diff --git a/docs/topics/spiders.rst b/docs/topics/spiders.rst index 2056664c7..a3e9f410f 100644 --- a/docs/topics/spiders.rst +++ b/docs/topics/spiders.rst @@ -294,6 +294,14 @@ The above example can also be written as follows:: def start_requests(self): yield scrapy.Request(f'http://www.example.com/categories/{self.category}') +If you are :ref:`running Scrapy from a script `, you can +specify spider arguments when calling +:class:`CrawlerProcess.crawl ` or +:class:`CrawlerRunner.crawl `:: + + process = CrawlerProcess() + process.crawl(MySpider, category="electronics") + Keep in mind that spider arguments are only strings. The spider will not do any parsing on its own. If you were to set the ``start_urls`` attribute from the command line, From c5b1ee810167266fcd259f263dbfc0fe0204761a Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Adri=C3=A1n=20Chaves?= Date: Tue, 11 May 2021 09:04:53 +0200 Subject: [PATCH 14/16] Make Twisted[http2] installation optional (#5113) Co-authored-by: Eugenio Lacuesta --- conftest.py | 9 +++++++++ docs/topics/settings.rst | 14 ++++++++----- setup.py | 3 +-- tests/test_downloader_handlers_http2.py | 26 ++++++++++++++++++++----- tests/test_http2_client_protocol.py | 13 ++++++++----- tox.ini | 11 +++++------ 6 files changed, 53 insertions(+), 23 deletions(-) diff --git a/conftest.py b/conftest.py index e4dd80de0..4931c5a79 100644 --- a/conftest.py +++ b/conftest.py @@ -1,6 +1,7 @@ from pathlib import Path import pytest +from twisted.web.http import H2_ENABLED from scrapy.utils.reactor import install_reactor @@ -25,6 +26,14 @@ for line in open('tests/ignores.txt'): if file_path and file_path[0] != '#': collect_ignore.append(file_path) +if not H2_ENABLED: + collect_ignore.extend( + ( + 'scrapy/core/downloader/handlers/http2.py', + *_py_files("scrapy/core/http2"), + ) + ) + @pytest.fixture() def chdir(tmpdir): diff --git a/docs/topics/settings.rst b/docs/topics/settings.rst index 2506497e2..0b290598f 100644 --- a/docs/topics/settings.rst +++ b/docs/topics/settings.rst @@ -680,12 +680,16 @@ handler (without replacement), place this in your ``settings.py``:: .. _http2: -The default HTTPS handler uses HTTP/1.1. To use HTTP/2 update -:setting:`DOWNLOAD_HANDLERS` as follows:: +The default HTTPS handler uses HTTP/1.1. To use HTTP/2: - DOWNLOAD_HANDLERS = { - 'https': 'scrapy.core.downloader.handlers.http2.H2DownloadHandler', - } +#. Install ``Twisted[http2]>=17.9.0`` to install the packages required to + enable HTTP/2 support in Twisted. + +#. Update :setting:`DOWNLOAD_HANDLERS` as follows:: + + DOWNLOAD_HANDLERS = { + 'https': 'scrapy.core.downloader.handlers.http2.H2DownloadHandler', + } .. warning:: diff --git a/setup.py b/setup.py index b1bb64575..ed2b6e347 100644 --- a/setup.py +++ b/setup.py @@ -19,7 +19,7 @@ def has_environment_marker_platform_impl_support(): install_requires = [ - 'Twisted[http2]>=17.9.0', + 'Twisted>=17.9.0', 'cryptography>=2.0', 'cssselect>=0.9.1', 'itemloaders>=1.0.1', @@ -31,7 +31,6 @@ install_requires = [ 'zope.interface>=4.1.3', 'protego>=0.1.15', 'itemadapter>=0.1.0', - 'h2>=3.0,<4.0', 'setuptools', ] extras_require = {} diff --git a/tests/test_downloader_handlers_http2.py b/tests/test_downloader_handlers_http2.py index 439778014..53bb4fe92 100644 --- a/tests/test_downloader_handlers_http2.py +++ b/tests/test_downloader_handlers_http2.py @@ -1,5 +1,5 @@ import json -from unittest import mock +from unittest import mock, skipIf from pytest import mark from testfixtures import LogCapture @@ -7,8 +7,8 @@ from twisted.internet import defer, error, reactor from twisted.trial import unittest from twisted.web import server from twisted.web.error import SchemeNotSupported +from twisted.web.http import H2_ENABLED -from scrapy.core.downloader.handlers.http2 import H2DownloadHandler from scrapy.http import Request from scrapy.spiders import Spider from scrapy.utils.misc import create_instance @@ -21,11 +21,17 @@ from tests.test_downloader_handlers import ( ) +@skipIf(not H2_ENABLED, "HTTP/2 support in Twisted is not enabled") class Https2TestCase(Https11TestCase): + scheme = 'https' - download_handler_cls = H2DownloadHandler HTTP2_DATALOSS_SKIP_REASON = "Content-Length mismatch raises InvalidBodyLengthError" + @classmethod + def setUpClass(cls): + from scrapy.core.downloader.handlers.http2 import H2DownloadHandler + cls.download_handler_cls = H2DownloadHandler + def test_protocol(self): request = Request(self.getURL("host"), method="GET") d = self.download_request(request, Spider("foo")) @@ -187,9 +193,14 @@ class Https2InvalidDNSPattern(Https2TestCase): super(Https2InvalidDNSPattern, self).setUp() +@skipIf(not H2_ENABLED, "HTTP/2 support in Twisted is not enabled") class Https2CustomCiphers(Https11CustomCiphers): scheme = 'https' - download_handler_cls = H2DownloadHandler + + @classmethod + def setUpClass(cls): + from scrapy.core.downloader.handlers.http2 import H2DownloadHandler + cls.download_handler_cls = H2DownloadHandler class Http2MockServerTestCase(Http11MockServerTestCase): @@ -201,6 +212,7 @@ class Http2MockServerTestCase(Http11MockServerTestCase): } +@skipIf(not H2_ENABLED, "HTTP/2 support in Twisted is not enabled") class Https2ProxyTestCase(Http11ProxyTestCase): # only used for HTTPS tests keyfile = 'keys/localhost.key' @@ -209,9 +221,13 @@ class Https2ProxyTestCase(Http11ProxyTestCase): scheme = 'https' host = u'127.0.0.1' - download_handler_cls = H2DownloadHandler expected_http_proxy_request_body = b'/' + @classmethod + def setUpClass(cls): + from scrapy.core.downloader.handlers.http2 import H2DownloadHandler + cls.download_handler_cls = H2DownloadHandler + def setUp(self): site = server.Site(UriResource(), timeout=None) self.port = reactor.listenSSL( diff --git a/tests/test_http2_client_protocol.py b/tests/test_http2_client_protocol.py index 8b2f6a11d..677ede92b 100644 --- a/tests/test_http2_client_protocol.py +++ b/tests/test_http2_client_protocol.py @@ -5,10 +5,9 @@ import re import shutil import string from ipaddress import IPv4Address -from unittest import mock +from unittest import mock, skipIf from urllib.parse import urlencode -from h2.exceptions import InvalidBodyLengthError from twisted.internet import reactor from twisted.internet.defer import CancelledError, Deferred, DeferredList, inlineCallbacks from twisted.internet.endpoints import SSL4ClientEndpoint, SSL4ServerEndpoint @@ -17,12 +16,10 @@ from twisted.internet.ssl import optionsForClientTLS, PrivateCertificate, Certif from twisted.python.failure import Failure from twisted.trial.unittest import TestCase from twisted.web.client import ResponseFailed, URI -from twisted.web.http import Request as TxRequest +from twisted.web.http import H2_ENABLED, Request as TxRequest from twisted.web.server import Site, NOT_DONE_YET from twisted.web.static import File -from scrapy.core.http2.protocol import H2ClientFactory, H2ClientProtocol -from scrapy.core.http2.stream import InactiveStreamClosed, InvalidHostname from scrapy.http import Request, Response, JsonRequest from scrapy.settings import Settings from scrapy.spiders import Spider @@ -173,6 +170,7 @@ def get_client_certificate(key_file, certificate_file) -> PrivateCertificate: return PrivateCertificate.loadPEM(pem) +@skipIf(not H2_ENABLED, "HTTP/2 support in Twisted is not enabled") class Https2ClientProtocolTestCase(TestCase): scheme = 'https' key_file = os.path.join(os.path.dirname(__file__), 'keys', 'localhost.key') @@ -220,6 +218,7 @@ class Https2ClientProtocolTestCase(TestCase): uri = URI.fromBytes(bytes(self.get_url('/'), 'utf-8')) self.conn_closed_deferred = Deferred() + from scrapy.core.http2.protocol import H2ClientFactory h2_client_factory = H2ClientFactory(uri, Settings(), self.conn_closed_deferred) client_endpoint = SSL4ClientEndpoint(reactor, self.hostname, self.port_number, client_options) self.client = yield client_endpoint.connect(h2_client_factory) @@ -426,6 +425,7 @@ class Https2ClientProtocolTestCase(TestCase): def assert_failure(failure: Failure): self.assertTrue(len(failure.value.reasons) > 0) + from h2.exceptions import InvalidBodyLengthError self.assertTrue(any( isinstance(error, InvalidBodyLengthError) for error in failure.value.reasons @@ -511,6 +511,7 @@ class Https2ClientProtocolTestCase(TestCase): def assert_inactive_stream(failure): self.assertIsNotNone(failure.check(ResponseFailed)) + from scrapy.core.http2.stream import InactiveStreamClosed self.assertTrue(any( isinstance(e, InactiveStreamClosed) for e in failure.value.reasons @@ -596,6 +597,7 @@ class Https2ClientProtocolTestCase(TestCase): request = Request(url) def assert_invalid_hostname(failure: Failure): + from scrapy.core.http2.stream import InvalidHostname self.assertIsNotNone(failure.check(InvalidHostname)) error_msg = str(failure.value) self.assertIn('localhost', error_msg) @@ -633,6 +635,7 @@ class Https2ClientProtocolTestCase(TestCase): def assert_timeout_error(failure: Failure): for err in failure.value.reasons: + from scrapy.core.http2.protocol import H2ClientProtocol if isinstance(err, TimeoutError): self.assertIn(f"Connection was IDLE for more than {H2ClientProtocol.IDLE_TIMEOUT}s", str(err)) break diff --git a/tox.ini b/tox.ini index 5b0606f8f..8167aff96 100644 --- a/tox.ini +++ b/tox.ini @@ -50,6 +50,8 @@ commands = basepython = python3 deps = {[testenv]deps} + # Twisted[http2] is required to import some files + Twisted[http2]>=17.9.0 pytest-flake8 commands = py.test --flake8 {posargs:docs scrapy tests} @@ -57,12 +59,7 @@ commands = [testenv:pylint] basepython = python3 deps = - {[testenv]deps} - # Optional dependencies - boto - reppy - robotexclusionrulesparser - # Test dependencies + {[testenv:extra-deps]deps} pylint commands = pylint conftest.py docs extras scrapy setup.py tests @@ -119,9 +116,11 @@ setenv = [testenv:extra-deps] deps = {[testenv]deps} + boto reppy robotexclusionrulesparser Pillow>=4.0.0 + Twisted[http2]>=17.9.0 [testenv:asyncio] commands = From ee682af3b06d48815dbdaa27c1177b94aaf679e1 Mon Sep 17 00:00:00 2001 From: Bhavesh <35660861+Bhavesh0327@users.noreply.github.com> Date: Wed, 12 May 2021 01:53:02 +0530 Subject: [PATCH 15/16] [Fix] Change the truncation limit of Proxy TunnelError from 32 to 1000 (#5007) * [Fix] Change the truncation limit oof Proxy TunnelError from 32 to 64 * [Fix] Change the truncation limit for Proxy tunnel error * [Fix] flake8 check * [Fix] formatting issues * [Remove] coverage report * [Fix] truncation error issue * [Fix] formatting issues * [Remove] coverage report --- scrapy/core/downloader/handlers/http11.py | 7 ++++--- 1 file changed, 4 insertions(+), 3 deletions(-) diff --git a/scrapy/core/downloader/handlers/http11.py b/scrapy/core/downloader/handlers/http11.py index 25cb3ec62..073f35891 100644 --- a/scrapy/core/downloader/handlers/http11.py +++ b/scrapy/core/downloader/handlers/http11.py @@ -98,8 +98,9 @@ class TunnelingTCP4ClientEndpoint(TCP4ClientEndpoint): with this endpoint comes from the pool and a CONNECT has already been issued for it. """ - - _responseMatcher = re.compile(br'HTTP/1\.. (?P\d{3})(?P.{,32})') + _truncatedLength = 1000 + _responseAnswer = r'HTTP/1\.. (?P\d{3})(?P.{,' + str(_truncatedLength) + r'})' + _responseMatcher = re.compile(_responseAnswer.encode()) def __init__(self, reactor, host, port, proxyConf, contextFactory, timeout=30, bindAddress=None): proxyHost, proxyPort, self._proxyAuthHeader = proxyConf @@ -144,7 +145,7 @@ class TunnelingTCP4ClientEndpoint(TCP4ClientEndpoint): extra = {'status': int(respm.group('status')), 'reason': respm.group('reason').strip()} else: - extra = rcvd_bytes[:32] + extra = rcvd_bytes[:self._truncatedLength] self._tunnelReadyDeferred.errback( TunnelError('Could not open CONNECT tunnel with proxy ' f'{self._host}:{self._port} [{extra!r}]') From 23cfdb058e80a18ee2b66e1b966355f3aca426d0 Mon Sep 17 00:00:00 2001 From: Vostretsov Nikita Date: Fri, 28 May 2021 09:45:06 +0000 Subject: [PATCH 16/16] Reducing amount of warnings during test run (#5162) * put flake8 options into separate file to remove pytest warnings * remove ResourceLeaked warning in pypy * suppress warnings from twisted * ignore deprecation warnings here * ignore deprecation warning in tests of deprecated methods * ignore deprecation warnings here * update test classes * don`t use deprecated method call * ignore deprecation warnings here * proper warning class * more selective ignoring * Revert "don`t use deprecated method call" This reverts commit 59216ab5603c4b47574382768614ef4c39d36747. --- .flake8 | 19 +++++++++++++++++++ .gitignore | 2 ++ conftest.py | 9 +++++---- pytest.ini | 19 ++----------------- tests/test_exporters.py | 12 ++++++++---- tests/test_feedexport.py | 4 ++-- tests/test_http_response.py | 12 ++++++++---- tests/test_item.py | 22 ++++++++++++---------- tests/test_utils_deprecate.py | 2 +- tests/test_utils_python.py | 9 +++++++-- 10 files changed, 66 insertions(+), 44 deletions(-) create mode 100644 .flake8 diff --git a/.flake8 b/.flake8 new file mode 100644 index 000000000..1c503fb0b --- /dev/null +++ b/.flake8 @@ -0,0 +1,19 @@ +[flake8] + +max-line-length = 119 +ignore = W503 + +exclude = +# Exclude files that are meant to provide top-level imports +# E402: Module level import not at top of file +# F401: Module imported but unused + scrapy/__init__.py E402 + scrapy/core/downloader/handlers/http.py F401 + scrapy/http/__init__.py F401 + scrapy/linkextractors/__init__.py E402 F401 + scrapy/selector/__init__.py F401 + scrapy/spiders/__init__.py E402 F401 + + # Issues pending a review: + scrapy/utils/url.py F403 F405 + tests/test_loader.py E741 diff --git a/.gitignore b/.gitignore index 795e2605e..d77d24624 100644 --- a/.gitignore +++ b/.gitignore @@ -14,6 +14,8 @@ htmlcov/ .coverage .pytest_cache/ .coverage.* +coverage.* +test-output.* .cache/ .mypy_cache/ /tests/keys/localhost.crt diff --git a/conftest.py b/conftest.py index 4931c5a79..05b4ccdad 100644 --- a/conftest.py +++ b/conftest.py @@ -21,10 +21,11 @@ collect_ignore = [ *_py_files("tests/CrawlerRunner"), ] -for line in open('tests/ignores.txt'): - file_path = line.strip() - if file_path and file_path[0] != '#': - collect_ignore.append(file_path) +with open('tests/ignores.txt') as reader: + for line in reader: + file_path = line.strip() + if file_path and file_path[0] != '#': + collect_ignore.append(file_path) if not H2_ENABLED: collect_ignore.extend( diff --git a/pytest.ini b/pytest.ini index 0aae09ff5..6de08c78d 100644 --- a/pytest.ini +++ b/pytest.ini @@ -20,20 +20,5 @@ addopts = --ignore=docs/utils markers = only_asyncio: marks tests as only enabled when --reactor=asyncio is passed -flake8-max-line-length = 119 -flake8-ignore = - W503 - - # Exclude files that are meant to provide top-level imports - # E402: Module level import not at top of file - # F401: Module imported but unused - scrapy/__init__.py E402 - scrapy/core/downloader/handlers/http.py F401 - scrapy/http/__init__.py F401 - scrapy/linkextractors/__init__.py E402 F401 - scrapy/selector/__init__.py F401 - scrapy/spiders/__init__.py E402 F401 - - # Issues pending a review: - scrapy/utils/url.py F403 F405 - tests/test_loader.py E741 +filterwarnings= + ignore::DeprecationWarning:twisted.web.test.test_webclient diff --git a/tests/test_exporters.py b/tests/test_exporters.py index ebc477e74..04bae31d3 100644 --- a/tests/test_exporters.py +++ b/tests/test_exporters.py @@ -6,12 +6,14 @@ import tempfile import unittest from io import BytesIO from datetime import datetime +from warnings import catch_warnings, filterwarnings import lxml.etree from itemadapter import ItemAdapter from scrapy.item import Item, Field from scrapy.utils.python import to_unicode +from scrapy.exceptions import ScrapyDeprecationWarning from scrapy.exporters import ( BaseItemExporter, PprintItemExporter, PickleItemExporter, CsvItemExporter, XmlItemExporter, JsonLinesItemExporter, JsonItemExporter, @@ -172,10 +174,12 @@ class PythonItemExporterTest(BaseItemExporterTest): self.assertEqual(type(exported['age'][0]['age'][0]), dict) def test_export_binary(self): - exporter = PythonItemExporter(binary=True) - value = self.item_class(name='John\xa3', age='22') - expected = {b'name': b'John\xc2\xa3', b'age': b'22'} - self.assertEqual(expected, exporter.export_item(value)) + with catch_warnings(): + filterwarnings('ignore', category=ScrapyDeprecationWarning) + exporter = PythonItemExporter(binary=True) + value = self.item_class(name='John\xa3', age='22') + expected = {b'name': b'John\xc2\xa3', b'age': b'22'} + self.assertEqual(expected, exporter.export_item(value)) def test_nonstring_types_item(self): item = self._get_nonstring_types_item() diff --git a/tests/test_feedexport.py b/tests/test_feedexport.py index d248824fc..df7ec4461 100644 --- a/tests/test_feedexport.py +++ b/tests/test_feedexport.py @@ -515,7 +515,7 @@ class FromCrawlerFileFeedStorage(FileFeedStorage, FromCrawlerMixin): class DummyBlockingFeedStorage(BlockingFeedStorage): - def __init__(self, uri): + def __init__(self, uri, *args, feed_options=None): self.path = file_uri_to_path(uri) def _store_in_thread(self, file): @@ -541,7 +541,7 @@ class LogOnStoreFileStorage: It can be used to make sure `store` method is invoked. """ - def __init__(self, uri): + def __init__(self, uri, feed_options=None): self.path = file_uri_to_path(uri) self.logger = getLogger() diff --git a/tests/test_http_response.py b/tests/test_http_response.py index f831ef5dc..04a594d03 100644 --- a/tests/test_http_response.py +++ b/tests/test_http_response.py @@ -1,6 +1,6 @@ import unittest from unittest import mock -from warnings import catch_warnings +from warnings import catch_warnings, filterwarnings from w3lib.encoding import resolve_encoding @@ -134,7 +134,9 @@ class BaseResponseTest(unittest.TestCase): assert isinstance(response.text, str) self._assert_response_encoding(response, encoding) self.assertEqual(response.body, body_bytes) - self.assertEqual(response.body_as_unicode(), body_unicode) + with catch_warnings(): + filterwarnings("ignore", category=ScrapyDeprecationWarning) + self.assertEqual(response.body_as_unicode(), body_unicode) self.assertEqual(response.text, body_unicode) def _assert_response_encoding(self, response, encoding): @@ -345,8 +347,10 @@ class TextResponseTest(BaseResponseTest): r1 = self.response_class('http://www.example.com', body=original_string, encoding='cp1251') # check body_as_unicode - self.assertTrue(isinstance(r1.body_as_unicode(), str)) - self.assertEqual(r1.body_as_unicode(), unicode_string) + with catch_warnings(): + filterwarnings("ignore", category=ScrapyDeprecationWarning) + self.assertTrue(isinstance(r1.body_as_unicode(), str)) + self.assertEqual(r1.body_as_unicode(), unicode_string) # check response.text self.assertTrue(isinstance(r1.text, str)) diff --git a/tests/test_item.py b/tests/test_item.py index 78d204e34..c94bb44af 100644 --- a/tests/test_item.py +++ b/tests/test_item.py @@ -1,6 +1,6 @@ import unittest from unittest import mock -from warnings import catch_warnings +from warnings import catch_warnings, filterwarnings from scrapy.exceptions import ScrapyDeprecationWarning from scrapy.item import ABCMeta, _BaseItem, BaseItem, DictItem, Field, Item, ItemMeta @@ -328,16 +328,18 @@ class BaseItemTest(unittest.TestCase): class SubclassedItem(Item): pass - self.assertTrue(isinstance(BaseItem(), BaseItem)) - self.assertTrue(isinstance(SubclassedBaseItem(), BaseItem)) - self.assertTrue(isinstance(Item(), BaseItem)) - self.assertTrue(isinstance(SubclassedItem(), BaseItem)) + with catch_warnings(): + filterwarnings("ignore", category=ScrapyDeprecationWarning) + self.assertTrue(isinstance(BaseItem(), BaseItem)) + self.assertTrue(isinstance(SubclassedBaseItem(), BaseItem)) + self.assertTrue(isinstance(Item(), BaseItem)) + self.assertTrue(isinstance(SubclassedItem(), BaseItem)) - # make sure internal checks using private _BaseItem class succeed - self.assertTrue(isinstance(BaseItem(), _BaseItem)) - self.assertTrue(isinstance(SubclassedBaseItem(), _BaseItem)) - self.assertTrue(isinstance(Item(), _BaseItem)) - self.assertTrue(isinstance(SubclassedItem(), _BaseItem)) + # make sure internal checks using private _BaseItem class succeed + self.assertTrue(isinstance(BaseItem(), _BaseItem)) + self.assertTrue(isinstance(SubclassedBaseItem(), _BaseItem)) + self.assertTrue(isinstance(Item(), _BaseItem)) + self.assertTrue(isinstance(SubclassedItem(), _BaseItem)) def test_deprecation_warning(self): """ diff --git a/tests/test_utils_deprecate.py b/tests/test_utils_deprecate.py index 35d35b45d..e47afa266 100644 --- a/tests/test_utils_deprecate.py +++ b/tests/test_utils_deprecate.py @@ -108,7 +108,7 @@ class WarnWhenSubclassedTest(unittest.TestCase): # ignore subclassing warnings with warnings.catch_warnings(): - warnings.simplefilter('ignore', ScrapyDeprecationWarning) + warnings.simplefilter('ignore', MyWarning) class UserClass(Deprecated): pass diff --git a/tests/test_utils_python.py b/tests/test_utils_python.py index 3115cc92f..4b3964154 100644 --- a/tests/test_utils_python.py +++ b/tests/test_utils_python.py @@ -5,8 +5,9 @@ import platform import unittest from datetime import datetime from itertools import count -from warnings import catch_warnings +from warnings import catch_warnings, filterwarnings +from scrapy.exceptions import ScrapyDeprecationWarning from scrapy.utils.python import ( memoizemethod_noargs, binary_is_text, equal_attributes, WeakKeyCache, get_func_args, to_bytes, to_unicode, @@ -160,7 +161,11 @@ class UtilsPythonTestCase(unittest.TestCase): pass _values = count() - wk = WeakKeyCache(lambda k: next(_values)) + + with catch_warnings(): + filterwarnings("ignore", category=ScrapyDeprecationWarning) + wk = WeakKeyCache(lambda k: next(_values)) + k = _Weakme() v = wk[k] self.assertEqual(v, wk[k])