mirror of https://github.com/scrapy/scrapy.git
More *_async APIs (#6997)
* Add ExecutionEngine.close_spider_async().
* Add Scraper.{open,close}_spider_async().
* Fix double engine stopping via Crawler.stop().
* Add ExecutionEngine.stop_async().
* Add versionadded to new async APIs.
* Add ExecutionEngine.close_async().
This commit is contained in:
parent
d27d0a4ed9
commit
a0b766f9e1
|
|
@ -14,7 +14,7 @@ from time import time
|
|||
from traceback import format_exc
|
||||
from typing import TYPE_CHECKING, Any
|
||||
|
||||
from twisted.internet.defer import CancelledError, Deferred, inlineCallbacks, succeed
|
||||
from twisted.internet.defer import CancelledError, Deferred, inlineCallbacks
|
||||
from twisted.python.failure import Failure
|
||||
|
||||
from scrapy import signals
|
||||
|
|
@ -27,10 +27,13 @@ from scrapy.exceptions import (
|
|||
ScrapyDeprecationWarning,
|
||||
)
|
||||
from scrapy.http import Request, Response
|
||||
from scrapy.utils.asyncio import AsyncioLoopingCall, create_looping_call
|
||||
from scrapy.utils.asyncio import (
|
||||
AsyncioLoopingCall,
|
||||
create_looping_call,
|
||||
is_asyncio_available,
|
||||
)
|
||||
from scrapy.utils.defer import (
|
||||
_schedule_coro,
|
||||
deferred_f_from_coro_f,
|
||||
deferred_from_coro,
|
||||
maybe_deferred_to_future,
|
||||
)
|
||||
|
|
@ -79,10 +82,10 @@ class _Slot:
|
|||
self.inprogress.remove(request)
|
||||
self._maybe_fire_closing()
|
||||
|
||||
def close(self) -> Deferred[None]:
|
||||
async def close(self) -> None:
|
||||
self.closing = Deferred()
|
||||
self._maybe_fire_closing()
|
||||
return self.closing
|
||||
await maybe_deferred_to_future(self.closing)
|
||||
|
||||
def _maybe_fire_closing(self) -> None:
|
||||
if self.closing is not None and not self.inprogress:
|
||||
|
|
@ -116,7 +119,9 @@ class ExecutionEngine:
|
|||
self.start_time: float | None = None
|
||||
self._start: AsyncIterator[Any] | None = None
|
||||
self._closewait: Deferred[None] | None = None
|
||||
self._start_request_processing_dfd: Deferred[None] | None = None
|
||||
self._start_request_processing_awaitable: (
|
||||
asyncio.Future[None] | Deferred[None] | None
|
||||
) = None
|
||||
downloader_cls: type[Downloader] = load_object(self.settings["DOWNLOADER"])
|
||||
try:
|
||||
self.scheduler_cls: type[BaseScheduler] = self._get_scheduler_class(
|
||||
|
|
@ -136,7 +141,8 @@ class ExecutionEngine:
|
|||
|
||||
self.scraper: Scraper = Scraper(crawler)
|
||||
except Exception:
|
||||
self.close()
|
||||
if hasattr(self, "downloader"):
|
||||
self.downloader.close()
|
||||
raise
|
||||
|
||||
def _get_scheduler_class(self, settings: BaseSettings) -> type[BaseScheduler]:
|
||||
|
|
@ -154,9 +160,15 @@ class ExecutionEngine:
|
|||
ScrapyDeprecationWarning,
|
||||
stacklevel=2,
|
||||
)
|
||||
return deferred_from_coro(self.start_async(_start_request_processing))
|
||||
return deferred_from_coro(
|
||||
self.start_async(_start_request_processing=_start_request_processing)
|
||||
)
|
||||
|
||||
async def start_async(self, _start_request_processing: bool = True) -> None:
|
||||
async def start_async(self, *, _start_request_processing: bool = True) -> None:
|
||||
"""Start the execution engine.
|
||||
|
||||
.. versionadded:: VERSION
|
||||
"""
|
||||
if self.running:
|
||||
raise RuntimeError("Engine already running")
|
||||
self.start_time = time()
|
||||
|
|
@ -167,46 +179,63 @@ class ExecutionEngine:
|
|||
self.running = True
|
||||
self._closewait = Deferred()
|
||||
if _start_request_processing:
|
||||
self._start_request_processing_dfd = self._start_request_processing()
|
||||
coro = self._start_request_processing()
|
||||
if is_asyncio_available():
|
||||
# not wrapping in a Deferred here to avoid https://github.com/twisted/twisted/issues/12470
|
||||
# (can happen when this is cancelled, e.g. in test_close_during_start_iteration())
|
||||
self._start_request_processing_awaitable = asyncio.ensure_future(coro)
|
||||
else:
|
||||
self._start_request_processing_awaitable = Deferred.fromCoroutine(coro)
|
||||
await maybe_deferred_to_future(self._closewait)
|
||||
|
||||
def stop(self) -> Deferred[None]:
|
||||
"""Gracefully stop the execution engine"""
|
||||
warnings.warn(
|
||||
"ExecutionEngine.stop() is deprecated, use stop_async() instead",
|
||||
ScrapyDeprecationWarning,
|
||||
stacklevel=2,
|
||||
)
|
||||
return deferred_from_coro(self.stop_async())
|
||||
|
||||
@deferred_f_from_coro_f
|
||||
async def _finish_stopping_engine(_: Any) -> None:
|
||||
await self.signals.send_catch_log_async(signal=signals.engine_stopped)
|
||||
if self._closewait:
|
||||
self._closewait.callback(None)
|
||||
async def stop_async(self) -> None:
|
||||
"""Gracefully stop the execution engine.
|
||||
|
||||
.. versionadded:: VERSION
|
||||
"""
|
||||
|
||||
if not self.running:
|
||||
raise RuntimeError("Engine not running")
|
||||
|
||||
self.running = False
|
||||
if self._start_request_processing_dfd is not None:
|
||||
self._start_request_processing_dfd.cancel()
|
||||
self._start_request_processing_dfd = None
|
||||
dfd = (
|
||||
self.close_spider(self.spider, reason="shutdown")
|
||||
if self.spider is not None
|
||||
else succeed(None)
|
||||
)
|
||||
return dfd.addBoth(_finish_stopping_engine)
|
||||
if self._start_request_processing_awaitable is not None:
|
||||
self._start_request_processing_awaitable.cancel()
|
||||
self._start_request_processing_awaitable = None
|
||||
if self.spider is not None:
|
||||
await self.close_spider_async(reason="shutdown")
|
||||
await self.signals.send_catch_log_async(signal=signals.engine_stopped)
|
||||
if self._closewait:
|
||||
self._closewait.callback(None)
|
||||
|
||||
def close(self) -> Deferred[None]:
|
||||
warnings.warn(
|
||||
"ExecutionEngine.close() is deprecated, use close_async() instead",
|
||||
ScrapyDeprecationWarning,
|
||||
stacklevel=2,
|
||||
)
|
||||
return deferred_from_coro(self.close_async())
|
||||
|
||||
async def close_async(self) -> None:
|
||||
"""
|
||||
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:
|
||||
return self.stop() # will also close spider and downloader
|
||||
if self.spider is not None:
|
||||
return self.close_spider(
|
||||
self.spider, reason="shutdown"
|
||||
await self.stop_async() # will also close spider and downloader
|
||||
elif self.spider is not None:
|
||||
await self.close_spider_async(
|
||||
reason="shutdown"
|
||||
) # will also close downloader
|
||||
if hasattr(self, "downloader"):
|
||||
elif hasattr(self, "downloader"):
|
||||
self.downloader.close()
|
||||
return succeed(None)
|
||||
|
||||
def pause(self) -> None:
|
||||
self.paused = True
|
||||
|
|
@ -242,7 +271,6 @@ class ExecutionEngine:
|
|||
)
|
||||
self._slot.nextcall.schedule()
|
||||
|
||||
@deferred_f_from_coro_f
|
||||
async def _start_request_processing(self) -> None:
|
||||
"""Starts consuming Spider.start() output and sending scheduled
|
||||
requests."""
|
||||
|
|
@ -266,13 +294,13 @@ class ExecutionEngine:
|
|||
return
|
||||
except Exception:
|
||||
# an error happened, log it and stop the engine
|
||||
self._start_request_processing_dfd = None
|
||||
self._start_request_processing_awaitable = None
|
||||
logger.error(
|
||||
"Error while processing requests from start()",
|
||||
exc_info=True,
|
||||
extra={"spider": self.spider},
|
||||
)
|
||||
await maybe_deferred_to_future(self.stop())
|
||||
await self.stop_async()
|
||||
|
||||
def _start_scheduled_requests(self) -> None:
|
||||
if self._slot is None or self._slot.closing is not None or self.paused:
|
||||
|
|
@ -409,7 +437,12 @@ class ExecutionEngine:
|
|||
return deferred_from_coro(self.download_async(request))
|
||||
|
||||
async def download_async(self, request: Request) -> Response:
|
||||
"""Asynchronous version of download() that returns a Response."""
|
||||
"""Return a coroutine which fires with a Response as result.
|
||||
|
||||
Only downloader middlewares are applied.
|
||||
|
||||
.. versionadded:: VERSION
|
||||
"""
|
||||
if self.spider is None:
|
||||
raise RuntimeError(f"No open spider to crawl: {request}")
|
||||
try:
|
||||
|
|
@ -481,7 +514,7 @@ class ExecutionEngine:
|
|||
self._start = await self.scraper.spidermw.process_start()
|
||||
if hasattr(scheduler, "open") and (d := scheduler.open(self.crawler.spider)):
|
||||
await maybe_deferred_to_future(d)
|
||||
await maybe_deferred_to_future(self.scraper.open_spider())
|
||||
await self.scraper.open_spider_async()
|
||||
assert self.crawler.stats
|
||||
self.crawler.stats.open_spider(self.crawler.spider)
|
||||
await self.signals.send_catch_log_async(
|
||||
|
|
@ -511,20 +544,33 @@ class ExecutionEngine:
|
|||
if self.spider_is_idle():
|
||||
ex = detected_ex.get(CloseSpider, CloseSpider(reason="finished"))
|
||||
assert isinstance(ex, CloseSpider) # typing
|
||||
self.close_spider(self.spider, reason=ex.reason)
|
||||
_schedule_coro(self.close_spider_async(reason=ex.reason))
|
||||
|
||||
def close_spider(self, spider: Spider, reason: str = "cancelled") -> Deferred[None]:
|
||||
warnings.warn(
|
||||
"ExecutionEngine.close_spider() is deprecated, use close_spider_async() instead",
|
||||
ScrapyDeprecationWarning,
|
||||
stacklevel=2,
|
||||
)
|
||||
return deferred_from_coro(self.close_spider_async(reason=reason))
|
||||
|
||||
async def close_spider_async(self, *, reason: str = "cancelled") -> None:
|
||||
"""Close (cancel) spider and clear all its outstanding requests.
|
||||
|
||||
.. versionadded:: VERSION
|
||||
"""
|
||||
if self.spider is None:
|
||||
raise RuntimeError("Spider not opened")
|
||||
|
||||
@inlineCallbacks
|
||||
def close_spider(
|
||||
self, spider: Spider, reason: str = "cancelled"
|
||||
) -> Generator[Deferred[Any], Any, None]:
|
||||
"""Close (cancel) spider and clear all its outstanding requests"""
|
||||
if self._slot is None:
|
||||
raise RuntimeError("Engine slot not assigned")
|
||||
|
||||
if self._slot.closing is not None:
|
||||
yield self._slot.closing
|
||||
await maybe_deferred_to_future(self._slot.closing)
|
||||
return
|
||||
|
||||
spider = self.spider
|
||||
|
||||
logger.info(
|
||||
"Closing spider (%(reason)s)", {"reason": reason}, extra={"spider": spider}
|
||||
)
|
||||
|
|
@ -533,7 +579,7 @@ class ExecutionEngine:
|
|||
logger.error(msg, exc_info=True, extra={"spider": spider}) # noqa: LOG014
|
||||
|
||||
try:
|
||||
yield self._slot.close()
|
||||
await self._slot.close()
|
||||
except Exception:
|
||||
log_failure("Slot close failure")
|
||||
|
||||
|
|
@ -543,19 +589,19 @@ class ExecutionEngine:
|
|||
log_failure("Downloader close failure")
|
||||
|
||||
try:
|
||||
yield self.scraper.close_spider()
|
||||
await self.scraper.close_spider_async()
|
||||
except Exception:
|
||||
log_failure("Scraper close failure")
|
||||
|
||||
if hasattr(self._slot.scheduler, "close"):
|
||||
try:
|
||||
if (d := self._slot.scheduler.close(reason)) is not None:
|
||||
yield d
|
||||
await maybe_deferred_to_future(d)
|
||||
except Exception:
|
||||
log_failure("Scheduler close failure")
|
||||
|
||||
try:
|
||||
yield self.signals.send_catch_log_deferred(
|
||||
await self.signals.send_catch_log_async(
|
||||
signal=signals.spider_closed,
|
||||
spider=spider,
|
||||
reason=reason,
|
||||
|
|
@ -580,6 +626,6 @@ class ExecutionEngine:
|
|||
|
||||
try:
|
||||
if (d := self._spider_closed_callback(spider)) is not None:
|
||||
yield d
|
||||
await maybe_deferred_to_future(d)
|
||||
except Exception:
|
||||
log_failure("Error running spider_closed_callback")
|
||||
|
|
|
|||
|
|
@ -24,6 +24,7 @@ from scrapy.http import Request, Response
|
|||
from scrapy.utils.asyncio import _parallel_asyncio, is_asyncio_available
|
||||
from scrapy.utils.defer import (
|
||||
_defer_sleep_async,
|
||||
_schedule_coro,
|
||||
aiter_errback,
|
||||
deferred_f_from_coro_f,
|
||||
deferred_from_coro,
|
||||
|
|
@ -130,16 +131,19 @@ class Scraper:
|
|||
assert crawler.logformatter
|
||||
self.logformatter: LogFormatter = crawler.logformatter
|
||||
|
||||
@deferred_f_from_coro_f
|
||||
async def open_spider(self, spider: Spider | None = None) -> None:
|
||||
"""Open the spider for scraping and allocate resources for it"""
|
||||
if spider is not None:
|
||||
warnings.warn(
|
||||
"Passing a 'spider' argument to Scraper.open_spider() is deprecated.",
|
||||
category=ScrapyDeprecationWarning,
|
||||
stacklevel=2,
|
||||
)
|
||||
def open_spider(self, spider: Spider | None = None) -> Deferred[None]:
|
||||
warnings.warn(
|
||||
"Scraper.open_spider() is deprecated, use open_spider_async() instead",
|
||||
ScrapyDeprecationWarning,
|
||||
stacklevel=2,
|
||||
)
|
||||
return deferred_from_coro(self.open_spider_async())
|
||||
|
||||
async def open_spider_async(self) -> None:
|
||||
"""Open the spider for scraping and allocate resources for it.
|
||||
|
||||
.. versionadded:: VERSION
|
||||
"""
|
||||
self.slot = Slot(self.crawler.settings.getint("SCRAPER_SLOT_MAX_ACTIVE_SIZE"))
|
||||
if not self.crawler.spider:
|
||||
raise RuntimeError(
|
||||
|
|
@ -152,24 +156,30 @@ class Scraper:
|
|||
else:
|
||||
await maybe_deferred_to_future(self.itemproc.open_spider())
|
||||
|
||||
def close_spider(self, spider: Spider | None = None) -> Deferred[list[None]]:
|
||||
"""Close the spider being scraped and release its resources"""
|
||||
if spider is not None:
|
||||
warnings.warn(
|
||||
"Passing a 'spider' argument to Scraper.close_spider() is deprecated.",
|
||||
category=ScrapyDeprecationWarning,
|
||||
stacklevel=2,
|
||||
)
|
||||
def close_spider(self, spider: Spider | None = None) -> Deferred[None]:
|
||||
warnings.warn(
|
||||
"Scraper.close_spider() is deprecated, use close_spider_async() instead",
|
||||
ScrapyDeprecationWarning,
|
||||
stacklevel=2,
|
||||
)
|
||||
return deferred_from_coro(self.close_spider_async())
|
||||
|
||||
async def close_spider_async(self) -> None:
|
||||
"""Close the spider being scraped and release its resources.
|
||||
|
||||
.. versionadded:: VERSION
|
||||
"""
|
||||
if self.slot is None:
|
||||
raise RuntimeError("Scraper slot not assigned")
|
||||
self.slot.closing = Deferred()
|
||||
if self._itemproc_needs_spider["close_spider"]:
|
||||
d = self.slot.closing.addCallback(self.itemproc.close_spider)
|
||||
else:
|
||||
d = self.slot.closing.addCallback(lambda _: self.itemproc.close_spider())
|
||||
self._check_if_closing()
|
||||
return d
|
||||
await maybe_deferred_to_future(self.slot.closing)
|
||||
if self._itemproc_needs_spider["close_spider"]:
|
||||
await maybe_deferred_to_future(
|
||||
self.itemproc.close_spider(self.crawler.spider)
|
||||
)
|
||||
else:
|
||||
await maybe_deferred_to_future(self.itemproc.close_spider())
|
||||
|
||||
def is_idle(self) -> bool:
|
||||
"""Return True if there isn't any more spiders to process"""
|
||||
|
|
@ -271,7 +281,10 @@ class Scraper:
|
|||
async def call_spider_async(
|
||||
self, result: Response | Failure, request: Request
|
||||
) -> Iterable[Any] | AsyncIterator[Any]:
|
||||
"""Call the request callback or errback with the response or failure."""
|
||||
"""Call the request callback or errback with the response or failure.
|
||||
|
||||
.. versionadded:: 2.13
|
||||
"""
|
||||
await _defer_sleep_async()
|
||||
assert self.crawler.spider
|
||||
if isinstance(result, Response):
|
||||
|
|
@ -315,8 +328,8 @@ class Scraper:
|
|||
exc = _failure.value
|
||||
if isinstance(exc, CloseSpider):
|
||||
assert self.crawler.engine is not None # typing
|
||||
self.crawler.engine.close_spider(
|
||||
self.crawler.spider, exc.reason or "cancelled"
|
||||
_schedule_coro(
|
||||
self.crawler.engine.close_spider_async(reason=exc.reason or "cancelled")
|
||||
)
|
||||
return
|
||||
logkws = self.logformatter.spider_error(
|
||||
|
|
@ -365,7 +378,10 @@ class Scraper:
|
|||
request: Request,
|
||||
response: Response | Failure,
|
||||
) -> None:
|
||||
"""Pass items/requests produced by a callback to ``_process_spidermw_output()`` in parallel."""
|
||||
"""Pass items/requests produced by a callback to ``_process_spidermw_output()`` in parallel.
|
||||
|
||||
.. versionadded:: 2.13
|
||||
"""
|
||||
it: Iterable[_T] | AsyncIterator[_T]
|
||||
if is_asyncio_available():
|
||||
if isinstance(result, AsyncIterator):
|
||||
|
|
@ -444,6 +460,8 @@ class Scraper:
|
|||
|
||||
*response* is the source of the item data. If the item does not come
|
||||
from response data, e.g. it was hard-coded, set it to ``None``.
|
||||
|
||||
.. versionadded:: VERSION
|
||||
"""
|
||||
assert self.slot is not None # typing
|
||||
assert self.crawler.spider is not None # typing
|
||||
|
|
|
|||
|
|
@ -164,7 +164,7 @@ class Crawler:
|
|||
except Exception:
|
||||
self.crawling = False
|
||||
if self.engine is not None:
|
||||
yield self.engine.close()
|
||||
yield deferred_from_coro(self.engine.close_async())
|
||||
raise
|
||||
|
||||
async def crawl_async(self, *args: Any, **kwargs: Any) -> None:
|
||||
|
|
@ -200,7 +200,7 @@ class Crawler:
|
|||
except Exception:
|
||||
self.crawling = False
|
||||
if self.engine is not None:
|
||||
await deferred_to_future(self.engine.close())
|
||||
await self.engine.close_async()
|
||||
raise
|
||||
|
||||
def _create_spider(self, *args: Any, **kwargs: Any) -> Spider:
|
||||
|
|
@ -216,7 +216,8 @@ class Crawler:
|
|||
if self.crawling:
|
||||
self.crawling = False
|
||||
assert self.engine
|
||||
yield self.engine.stop()
|
||||
if self.engine.running:
|
||||
yield deferred_from_coro(self.engine.stop_async())
|
||||
|
||||
async def stop_async(self) -> None:
|
||||
"""Start a graceful stop of the crawler and complete when the crawler is stopped.
|
||||
|
|
|
|||
|
|
@ -18,6 +18,7 @@ from scrapy.utils.asyncio import (
|
|||
call_later,
|
||||
create_looping_call,
|
||||
)
|
||||
from scrapy.utils.defer import _schedule_coro
|
||||
|
||||
if TYPE_CHECKING:
|
||||
from twisted.internet.task import LoopingCall
|
||||
|
|
@ -86,38 +87,31 @@ class CloseSpider:
|
|||
def error_count(self, failure: Failure, response: Response, spider: Spider) -> None:
|
||||
self.counter["errorcount"] += 1
|
||||
if self.counter["errorcount"] == self.close_on["errorcount"]:
|
||||
assert self.crawler.engine
|
||||
self.crawler.engine.close_spider(spider, "closespider_errorcount")
|
||||
self._close_spider("closespider_errorcount")
|
||||
|
||||
def page_count(self, response: Response, request: Request, spider: Spider) -> None:
|
||||
self.counter["pagecount"] += 1
|
||||
self.counter["pagecount_since_last_item"] += 1
|
||||
if self.counter["pagecount"] == self.close_on["pagecount"]:
|
||||
assert self.crawler.engine
|
||||
self.crawler.engine.close_spider(spider, "closespider_pagecount")
|
||||
self._close_spider("closespider_pagecount")
|
||||
return
|
||||
if self.close_on["pagecount_no_item"] and (
|
||||
self.counter["pagecount_since_last_item"]
|
||||
>= self.close_on["pagecount_no_item"]
|
||||
):
|
||||
assert self.crawler.engine
|
||||
self.crawler.engine.close_spider(spider, "closespider_pagecount_no_item")
|
||||
self._close_spider("closespider_pagecount_no_item")
|
||||
|
||||
def spider_opened(self, spider: Spider) -> None:
|
||||
assert self.crawler.engine
|
||||
self.task = call_later(
|
||||
self.close_on["timeout"],
|
||||
self.crawler.engine.close_spider,
|
||||
spider,
|
||||
"closespider_timeout",
|
||||
self.close_on["timeout"], self._close_spider, "closespider_timeout"
|
||||
)
|
||||
|
||||
def item_scraped(self, item: Any, spider: Spider) -> None:
|
||||
self.counter["itemcount"] += 1
|
||||
self.counter["pagecount_since_last_item"] = 0
|
||||
if self.counter["itemcount"] == self.close_on["itemcount"]:
|
||||
assert self.crawler.engine
|
||||
self.crawler.engine.close_spider(spider, "closespider_itemcount")
|
||||
self._close_spider("closespider_itemcount")
|
||||
|
||||
def spider_closed(self, spider: Spider) -> None:
|
||||
if self.task:
|
||||
|
|
@ -130,7 +124,7 @@ class CloseSpider:
|
|||
self.task_no_item = None
|
||||
|
||||
def spider_opened_no_item(self, spider: Spider) -> None:
|
||||
self.task_no_item = create_looping_call(self._count_items_produced, spider)
|
||||
self.task_no_item = create_looping_call(self._count_items_produced)
|
||||
self.task_no_item.start(self.timeout_no_item, now=False)
|
||||
|
||||
logger.info(
|
||||
|
|
@ -141,7 +135,7 @@ class CloseSpider:
|
|||
def item_scraped_no_item(self, item: Any, spider: Spider) -> None:
|
||||
self.items_in_period += 1
|
||||
|
||||
def _count_items_produced(self, spider: Spider) -> None:
|
||||
def _count_items_produced(self) -> None:
|
||||
if self.items_in_period >= 1:
|
||||
self.items_in_period = 0
|
||||
else:
|
||||
|
|
@ -149,5 +143,8 @@ class CloseSpider:
|
|||
f"Closing spider since no items were produced in the last "
|
||||
f"{self.timeout_no_item} seconds."
|
||||
)
|
||||
assert self.crawler.engine
|
||||
self.crawler.engine.close_spider(spider, "closespider_timeout_no_item")
|
||||
self._close_spider("closespider_timeout_no_item")
|
||||
|
||||
def _close_spider(self, reason: str) -> None:
|
||||
assert self.crawler.engine
|
||||
_schedule_coro(self.crawler.engine.close_spider_async(reason=reason))
|
||||
|
|
|
|||
|
|
@ -17,6 +17,7 @@ from scrapy import signals
|
|||
from scrapy.exceptions import NotConfigured
|
||||
from scrapy.mail import MailSender
|
||||
from scrapy.utils.asyncio import AsyncioLoopingCall, create_looping_call
|
||||
from scrapy.utils.defer import _schedule_coro
|
||||
from scrapy.utils.engine import get_engine_status
|
||||
|
||||
if TYPE_CHECKING:
|
||||
|
|
@ -110,8 +111,8 @@ class MemoryUsage:
|
|||
self.crawler.stats.set_value("memusage/limit_notified", 1)
|
||||
|
||||
if self.crawler.engine.spider is not None:
|
||||
self.crawler.engine.close_spider(
|
||||
self.crawler.engine.spider, "memusage_exceeded"
|
||||
_schedule_coro(
|
||||
self.crawler.engine.close_spider_async(reason="memusage_exceeded")
|
||||
)
|
||||
else:
|
||||
self.crawler.stop()
|
||||
|
|
|
|||
|
|
@ -25,7 +25,7 @@ from scrapy.spiders import Spider
|
|||
from scrapy.utils.conf import get_config
|
||||
from scrapy.utils.console import DEFAULT_PYTHON_SHELLS, start_python_console
|
||||
from scrapy.utils.datatypes import SequenceExclude
|
||||
from scrapy.utils.defer import deferred_f_from_coro_f
|
||||
from scrapy.utils.defer import _schedule_coro, deferred_f_from_coro_f
|
||||
from scrapy.utils.misc import load_object
|
||||
from scrapy.utils.reactor import is_asyncio_reactor_installed, set_asyncio_event_loop
|
||||
from scrapy.utils.response import open_in_browser
|
||||
|
|
@ -127,7 +127,7 @@ class Shell:
|
|||
self.crawler.spider = spider
|
||||
assert self.crawler.engine
|
||||
await self.crawler.engine.open_spider_async(close_if_idle=False)
|
||||
self.crawler.engine._start_request_processing()
|
||||
_schedule_coro(self.crawler.engine._start_request_processing())
|
||||
self.spider = spider
|
||||
|
||||
def fetch(
|
||||
|
|
|
|||
|
|
@ -77,6 +77,8 @@ class SignalManager:
|
|||
|
||||
The keyword arguments are passed to the signal handlers (connected
|
||||
through the :meth:`connect` method).
|
||||
|
||||
.. versionadded:: VERSION
|
||||
"""
|
||||
kwargs.setdefault("sender", self.sender)
|
||||
return await _signal.send_catch_log_async(signal, **kwargs)
|
||||
|
|
|
|||
|
|
@ -120,6 +120,8 @@ async def send_catch_log_async(
|
|||
<signal-deferred>`.
|
||||
|
||||
Returns a coroutine that completes once all signal handlers have finished.
|
||||
|
||||
.. versionadded:: VERSION
|
||||
"""
|
||||
return await maybe_deferred_to_future(
|
||||
send_catch_log_deferred(signal, sender, *arguments, **named)
|
||||
|
|
|
|||
|
|
@ -33,7 +33,7 @@ class TestManagerBase:
|
|||
crawler.engine = crawler._create_engine()
|
||||
await crawler.engine.open_spider_async()
|
||||
yield mwman
|
||||
await maybe_deferred_to_future(crawler.engine.close_spider(crawler.spider))
|
||||
await crawler.engine.close_spider_async()
|
||||
|
||||
@staticmethod
|
||||
async def _download(
|
||||
|
|
|
|||
|
|
@ -32,7 +32,6 @@ from scrapy.utils.defer import (
|
|||
_schedule_coro,
|
||||
deferred_f_from_coro_f,
|
||||
deferred_from_coro,
|
||||
deferred_to_future,
|
||||
maybe_deferred_to_future,
|
||||
)
|
||||
from scrapy.utils.signal import disconnect_all
|
||||
|
|
@ -417,10 +416,10 @@ class TestEngine(TestEngineBase):
|
|||
"reason": "custom_reason",
|
||||
} == run.signals_caught[signals.spider_closed]
|
||||
|
||||
@inlineCallbacks
|
||||
def test_close_downloader(self):
|
||||
@deferred_f_from_coro_f
|
||||
async def test_close_downloader(self):
|
||||
e = ExecutionEngine(get_crawler(MySpider), lambda _: None)
|
||||
yield e.close()
|
||||
await e.close_async()
|
||||
|
||||
def test_close_without_downloader(self):
|
||||
class CustomException(Exception):
|
||||
|
|
@ -444,7 +443,7 @@ class TestEngine(TestEngineBase):
|
|||
_schedule_coro(e.start_async())
|
||||
with pytest.raises(RuntimeError, match="Engine already running"):
|
||||
yield deferred_from_coro(e.start_async())
|
||||
yield e.stop()
|
||||
yield deferred_from_coro(e.stop_async())
|
||||
|
||||
@pytest.mark.only_asyncio
|
||||
@deferred_f_from_coro_f
|
||||
|
|
@ -455,7 +454,7 @@ class TestEngine(TestEngineBase):
|
|||
await e.open_spider_async()
|
||||
with pytest.raises(RuntimeError, match="Engine already running"):
|
||||
await asyncio.gather(e.start_async(), e.start_async())
|
||||
await deferred_to_future(e.stop())
|
||||
await e.stop_async()
|
||||
|
||||
@inlineCallbacks
|
||||
def test_start_request_processing_exception(self):
|
||||
|
|
@ -635,7 +634,7 @@ def test_request_scheduled_signal(caplog):
|
|||
|
||||
|
||||
class TestEngineCloseSpider:
|
||||
"""Tests for exception handling coverage during close_spider()."""
|
||||
"""Tests for exception handling coverage during close_spider_async()."""
|
||||
|
||||
@pytest.fixture
|
||||
def crawler(self) -> Crawler:
|
||||
|
|
@ -648,9 +647,14 @@ class TestEngineCloseSpider:
|
|||
engine = ExecutionEngine(crawler, lambda _: None)
|
||||
await engine.open_spider_async()
|
||||
engine._slot = None
|
||||
assert crawler.spider
|
||||
with pytest.raises(RuntimeError, match="Engine slot not assigned"):
|
||||
await maybe_deferred_to_future(engine.close_spider(crawler.spider))
|
||||
await engine.close_spider_async()
|
||||
|
||||
@deferred_f_from_coro_f
|
||||
async def test_no_spider(self, crawler: Crawler) -> None:
|
||||
engine = ExecutionEngine(crawler, lambda _: None)
|
||||
with pytest.raises(RuntimeError, match="Spider not opened"):
|
||||
await engine.close_spider_async()
|
||||
|
||||
@deferred_f_from_coro_f
|
||||
async def test_exception_slot(
|
||||
|
|
@ -660,8 +664,7 @@ class TestEngineCloseSpider:
|
|||
await engine.open_spider_async()
|
||||
assert engine._slot
|
||||
del engine._slot.heartbeat
|
||||
assert crawler.spider
|
||||
await maybe_deferred_to_future(engine.close_spider(crawler.spider))
|
||||
await engine.close_spider_async()
|
||||
assert "Slot close failure" in caplog.text
|
||||
|
||||
@deferred_f_from_coro_f
|
||||
|
|
@ -671,8 +674,7 @@ class TestEngineCloseSpider:
|
|||
engine = ExecutionEngine(crawler, lambda _: None)
|
||||
await engine.open_spider_async()
|
||||
del engine.downloader.slots
|
||||
assert crawler.spider
|
||||
await maybe_deferred_to_future(engine.close_spider(crawler.spider))
|
||||
await engine.close_spider_async()
|
||||
assert "Downloader close failure" in caplog.text
|
||||
|
||||
@deferred_f_from_coro_f
|
||||
|
|
@ -682,8 +684,7 @@ class TestEngineCloseSpider:
|
|||
engine = ExecutionEngine(crawler, lambda _: None)
|
||||
await engine.open_spider_async()
|
||||
engine.scraper.slot = None
|
||||
assert crawler.spider
|
||||
await maybe_deferred_to_future(engine.close_spider(crawler.spider))
|
||||
await engine.close_spider_async()
|
||||
assert "Scraper close failure" in caplog.text
|
||||
|
||||
@deferred_f_from_coro_f
|
||||
|
|
@ -694,8 +695,7 @@ class TestEngineCloseSpider:
|
|||
await engine.open_spider_async()
|
||||
assert engine._slot
|
||||
del cast("Scheduler", engine._slot.scheduler).dqs
|
||||
assert crawler.spider
|
||||
await maybe_deferred_to_future(engine.close_spider(crawler.spider))
|
||||
await engine.close_spider_async()
|
||||
assert "Scheduler close failure" in caplog.text
|
||||
|
||||
@deferred_f_from_coro_f
|
||||
|
|
@ -705,8 +705,7 @@ class TestEngineCloseSpider:
|
|||
engine = ExecutionEngine(crawler, lambda _: None)
|
||||
await engine.open_spider_async()
|
||||
del engine.signals
|
||||
assert crawler.spider
|
||||
await maybe_deferred_to_future(engine.close_spider(crawler.spider))
|
||||
await engine.close_spider_async()
|
||||
assert "Error while sending spider_close signal" in caplog.text
|
||||
|
||||
@deferred_f_from_coro_f
|
||||
|
|
@ -716,8 +715,7 @@ class TestEngineCloseSpider:
|
|||
engine = ExecutionEngine(crawler, lambda _: None)
|
||||
await engine.open_spider_async()
|
||||
del cast("MemoryStatsCollector", crawler.stats).spider_stats
|
||||
assert crawler.spider
|
||||
await maybe_deferred_to_future(engine.close_spider(crawler.spider))
|
||||
await engine.close_spider_async()
|
||||
assert "Stats close failure" in caplog.text
|
||||
|
||||
@deferred_f_from_coro_f
|
||||
|
|
@ -726,6 +724,5 @@ class TestEngineCloseSpider:
|
|||
) -> None:
|
||||
engine = ExecutionEngine(crawler, lambda _: defer.fail(ValueError()))
|
||||
await engine.open_spider_async()
|
||||
assert crawler.spider
|
||||
await maybe_deferred_to_future(engine.close_spider(crawler.spider))
|
||||
await engine.close_spider_async()
|
||||
assert "Error running spider_closed_callback" in caplog.text
|
||||
|
|
|
|||
|
|
@ -4,7 +4,6 @@ from collections import deque
|
|||
from logging import ERROR
|
||||
from typing import TYPE_CHECKING
|
||||
|
||||
from testfixtures import LogCapture
|
||||
from twisted.internet.defer import Deferred
|
||||
|
||||
from scrapy import Request, Spider, signals
|
||||
|
|
@ -14,6 +13,8 @@ from tests.mockserver.http import MockServer
|
|||
from tests.test_scheduler import MemoryScheduler
|
||||
|
||||
if TYPE_CHECKING:
|
||||
import pytest
|
||||
|
||||
from scrapy.http import Response
|
||||
|
||||
|
||||
|
|
@ -86,13 +87,15 @@ class TestMain:
|
|||
assert actual_urls == expected_urls, f"{actual_urls=} != {expected_urls=}"
|
||||
|
||||
@deferred_f_from_coro_f
|
||||
async def test_close_during_start_iteration(self):
|
||||
async def test_close_during_start_iteration(
|
||||
self, caplog: pytest.LogCaptureFixture
|
||||
) -> None:
|
||||
class TestSpider(Spider):
|
||||
name = "test"
|
||||
|
||||
async def start(self):
|
||||
assert self.crawler.engine is not None
|
||||
await maybe_deferred_to_future(self.crawler.engine.close())
|
||||
await self.crawler.engine.close_async()
|
||||
yield Request("data:,a")
|
||||
|
||||
def parse(self, response):
|
||||
|
|
@ -107,15 +110,14 @@ class TestMain:
|
|||
crawler = get_crawler(TestSpider, settings_dict=settings)
|
||||
crawler.signals.connect(track_url, signals.request_reached_downloader)
|
||||
|
||||
with LogCapture(level=ERROR) as log:
|
||||
caplog.clear()
|
||||
with caplog.at_level(ERROR):
|
||||
await maybe_deferred_to_future(crawler.crawl())
|
||||
|
||||
assert len(log.records) == 1
|
||||
assert log.records[0].msg == "Error running spider_closed_callback"
|
||||
finish_reason = crawler.stats.get_value("finish_reason")
|
||||
assert finish_reason == "shutdown", f"{finish_reason=}"
|
||||
expected_urls = []
|
||||
assert actual_urls == expected_urls, f"{actual_urls=} != {expected_urls=}"
|
||||
assert not caplog.records
|
||||
assert crawler.stats
|
||||
assert crawler.stats.get_value("finish_reason") == "shutdown"
|
||||
assert not actual_urls
|
||||
|
||||
|
||||
class TestRequestSendOrder:
|
||||
|
|
|
|||
Loading…
Reference in New Issue