Source code for silkworm.engine

"""Asynchronous crawl scheduler, lifecycle manager, and callback dispatcher."""

from __future__ import annotations

import asyncio
import hashlib
import reprlib
import sys
import time
from collections import Counter, deque
from collections.abc import Awaitable, Callable, Iterable
from contextlib import suppress
from contextvars import ContextVar
from dataclasses import dataclass
from datetime import timedelta
from functools import partial
from itertools import count
from typing import TYPE_CHECKING, TypedDict, cast

try:  # resource is POSIX-only
    import resource
except ImportError:  # pragma: no cover - platform dependent
    resource = None

from ._domains import DomainSlots, url_host
from ._jobs import JobState
from ._metrics import MetricsServer, render_metrics
from ._scope import CallbackContractError, CrawlScope, enter_scope, invoke_callback
from ._stats import CrawlResult, CrawlStats, freeze_result
from ._timeouts import to_seconds
from ._types import JSONLike, JSONValue
from ._urls import host_in_domains, normalize_domains, request_fingerprint
from ._validation import require_positive_int
from .exceptions import (
    CloseSpider,
    CrawlFailedError,
    DropItem,
    IgnoreRequest,
    SilkwormError,
    SpiderError,
)
from .http import (
    DEFAULT_EMULATION,
    DEFAULT_MAX_RESPONSE_SIZE_BYTES,
    DEFAULT_REQUEST_TIMEOUT,
    HttpClient,
)
from .logging import Logger, LogLevel, complete_logs, get_logger, log_at_level
from .request import Callback, Request
from .response import HTMLResponse, Response

if TYPE_CHECKING:
    import os

    from wreq import Emulation, Profile

    from .http import FetchClient
    from .httpcache import HttpCache
    from .middlewares import (
        ExceptionMiddleware,
        RequestMiddleware,
        ResponseMiddleware,
    )
    from .pipelines import ItemPipeline
    from .spiders import Spider


[docs] @dataclass(slots=True) class EngineLogger: """ Customizable engine event logger. Subclass this when you need to redact or reshape selected engine log events. Set an event level to ``None`` to suppress that event. """ fetched_response_level: LogLevel = "INFO" fetching_request_level: LogLevel = "DEBUG" item_pipeline_level: LogLevel = "DEBUG" retry_request_level: LogLevel = "DEBUG" include_request_url: bool = True
[docs] def fetching_request( self, logger: Logger, request: Request, spider: Spider, ) -> None: """Log that ``request`` is about to be sent.""" context: dict[str, object] = { "method": request.method, "callback": getattr(request.callback, "__name__", None), "spider": spider.name, } if self.include_request_url: context["url"] = request.url log_at_level(logger, self.fetching_request_level, "Fetching request", **context)
[docs] def fetched_response( self, logger: Logger, request: Request, response: Response, spider: Spider, ) -> None: """Log a completed response with status and optional request URL.""" context: dict[str, object] = { "status": response.status, "spider": spider.name, } if self.include_request_url: context["url"] = request.url log_at_level(logger, self.fetched_response_level, "Fetched response", **context)
[docs] def retrying_request( self, logger: Logger, request: Request, spider: Spider, *, source: str, ) -> None: """Log that a middleware or response requested another attempt.""" context: dict[str, object] = {"source": source, "spider": spider.name} if self.include_request_url: context["url"] = request.url log_at_level( logger, self.retry_request_level, "Retrying request", **context, )
[docs] def running_item_pipeline( self, logger: Logger, pipeline: ItemPipeline, spider: Spider, ) -> None: """Log item dispatch using the pipeline's effective log level.""" log_level = cast( "LogLevel", getattr(pipeline, "log_level", self.item_pipeline_level) ) log_at_level( logger, log_level, "Running item pipeline", pipeline=pipeline.__class__.__name__, spider=spider.name, )
_SAFE_REPR = reprlib.Repr() _SAFE_REPR.maxstring = 120 _SAFE_REPR.maxother = 120 _SAFE_REPR.maxlist = 8 _SAFE_REPR.maxdict = 8 _SAFE_REPR.maxset = 8 _SAFE_REPR.maxtuple = 8 type DedupKey = Callable[[Request], str] type PrioritizedRequest = tuple[int, int, Request] type LifecycleCloser = tuple[str, Callable[[], Awaitable[object]]] # Index of the engine worker whose callback (or a task it spawned) is running; # ``None`` outside workers, e.g. while ``start_requests()`` seeds the queue. _CURRENT_WORKER: ContextVar[int | None] = ContextVar( "silkworm_engine_worker", default=None, ) # ``items_dropped_by_reason`` label for items discarded because ``max_items`` # was reached; excluded from the item drop rate used by the failure policy. MAX_ITEMS_DROP_REASON = "max_items" class _CrawlClosed(BaseException): """Aborts ``start_requests()`` once the crawl is stopping. Derives from ``BaseException`` so ``except Exception`` blocks in spider code cannot swallow it. """
[docs] def default_dedup_key(req: Request) -> str: """Return the engine's default deduplication key for ``req``. This is :func:`~silkworm.request_fingerprint`: the HTTP method, the canonical URL (normalized case, default port, sorted query, no fragment) with :attr:`~silkworm.Request.params` merged in, and the request body. Headers and metadata do not affect it. """ return request_fingerprint(req)
[docs] class EngineOptions(TypedDict, total=False): """Keyword options for :class:`Engine`, also accepted by every runner. ``run_spider(MySpider, concurrency=32, request_timeout=10)`` forwards these to ``Engine``; omitted keys use the ``Engine`` defaults. Supply ``http_client`` to inject a compatible client; its concurrency then controls worker count and default queue capacity. """ concurrency: int max_pending_requests: int | None emulation: Emulation | Profile | None request_timeout: float | timedelta | None html_max_size_bytes: int max_response_size_bytes: int | None request_middlewares: Iterable[RequestMiddleware] | None response_middlewares: Iterable[ResponseMiddleware] | None item_pipelines: Iterable[ItemPipeline] | None item_batch_size: int item_batch_wait: float log_stats_interval: float | None keep_alive: bool http_client: FetchClient | None engine_logger: EngineLogger | None dedup_key: DedupKey | None concurrency_per_domain: int | None max_depth: int | None max_requests: int | None max_items: int | None max_errors: int | None max_duration: float | timedelta | None max_error_rate: float | None min_items: int | None max_item_drop_rate: float | None job_dir: str | os.PathLike[str] | None http_cache: HttpCache | None metrics_port: int | None metrics_host: str
def _require_optional_positive(value: int | None, name: str) -> None: if value is not None: require_positive_int(value, name) def _require_optional_rate(value: float | None, name: str) -> None: if value is not None and not 0.0 <= value <= 1.0: msg = f"{name} must be between 0.0 and 1.0" raise ValueError(msg) def _meta_depth(request: Request) -> int: depth = request.meta.get("depth", 0) return depth if isinstance(depth, int) and not isinstance(depth, bool) else 0
[docs] class Engine: """Coordinate request scheduling, HTTP I/O, callbacks, and item pipelines. Args: spider: Spider instance to execute. concurrency: Maximum simultaneous HTTP requests for the default client. max_pending_requests: Queue capacity used for backpressure. Defaults to ten times the effective HTTP client concurrency. ``start_requests()`` waits while the queue is full; callbacks wait only when they are the sole producer, otherwise they enqueue past the bound so workers keep crawling and can never deadlock. emulation: Browser profile used by the default ``wreq`` client; pass ``None`` to disable impersonation. request_timeout: Default per-request timeout for the default client: 60 seconds (:data:`~silkworm.http.DEFAULT_REQUEST_TIMEOUT`) unless given; ``None`` disables it. ``Request.timeout`` overrides it per request. Ignored when ``http_client`` is supplied. html_max_size_bytes: Maximum document size parsed by HTML responses. max_response_size_bytes: Largest response body the default client downloads (``None`` for no limit); larger bodies fail with :class:`~silkworm.exceptions.ResponseTooLargeError`. request_middlewares: Request processors applied in list order. response_middlewares: Response processors applied in list order. item_pipelines: Item processors applied in list order, passing each returned value to the next pipeline. item_batch_size: Number of emitted items processed together. The default ``1`` preserves immediate per-item processing. item_batch_wait: Maximum seconds to wait for a partial batch when ``item_batch_size`` is greater than one. log_stats_interval: Seconds between statistics messages, or ``None`` to disable periodic summaries. keep_alive: Request connection reuse from the default client when supported by ``wreq``. http_client: Preconfigured client replacing the default client. engine_logger: Event logger customization. dedup_key: Function mapping a request to its deduplication key. Defaults to :func:`default_dedup_key`. concurrency_per_domain: Maximum simultaneous fetches per host, or ``None`` for no per-host limit. max_depth: Drop requests more than this many links away from a start request (start requests have depth ``0``). max_requests: Stop after sending this many requests. max_items: Stop after this many items passed every pipeline; later items are dropped. max_errors: Stop after this many unrecovered failures. max_duration: Stop after this much wall-clock time. max_error_rate: Fail the crawl when ``errors / requests_sent`` exceeds this fraction. min_items: Fail the crawl when fewer items were scraped. max_item_drop_rate: Fail the crawl when the share of items dropped by pipelines exceeds this fraction. job_dir: Directory persisting the seen-set and unfinished requests so an interrupted crawl resumes where it stopped. http_cache: Serve and store responses through this on-disk cache. metrics_port: Serve Prometheus metrics at ``/metrics`` on this port while crawling (``0`` picks a free port). metrics_host: Interface for the metrics server. Requests with :attr:`~silkworm.Request.dont_filter` bypass deduplication and off-site filtering. Higher request priorities are dequeued before lower ones, while insertion order breaks ties. The stop limits end the crawl gracefully: pending requests are discarded (or kept in ``job_dir``), in-flight requests finish, and :meth:`run` reports the limit as the close reason. Failure-policy violations make :meth:`run` raise :class:`~silkworm.exceptions.CrawlFailedError`; they are not evaluated for crawls stopped with :meth:`stop`. """
[docs] def __init__( self, spider: Spider, *, concurrency: int = 16, max_pending_requests: int | None = None, emulation: Emulation | Profile | None = DEFAULT_EMULATION, request_timeout: float | timedelta | None = DEFAULT_REQUEST_TIMEOUT, html_max_size_bytes: int = 5_000_000, max_response_size_bytes: int | None = DEFAULT_MAX_RESPONSE_SIZE_BYTES, request_middlewares: Iterable[RequestMiddleware] | None = None, response_middlewares: Iterable[ResponseMiddleware] | None = None, item_pipelines: Iterable[ItemPipeline] | None = None, item_batch_size: int = 1, item_batch_wait: float = 0.05, log_stats_interval: float | None = None, keep_alive: bool = False, http_client: FetchClient | None = None, engine_logger: EngineLogger | None = None, dedup_key: DedupKey | None = None, concurrency_per_domain: int | None = None, max_depth: int | None = None, max_requests: int | None = None, max_items: int | None = None, max_errors: int | None = None, max_duration: float | timedelta | None = None, max_error_rate: float | None = None, min_items: int | None = None, max_item_drop_rate: float | None = None, job_dir: str | os.PathLike[str] | None = None, http_cache: HttpCache | None = None, metrics_port: int | None = None, metrics_host: str = "127.0.0.1", ) -> None: require_positive_int(concurrency, "concurrency") require_positive_int(item_batch_size, "item_batch_size") if item_batch_wait <= 0: raise ValueError("item_batch_wait must be positive") for name, value in ( ("max_requests", max_requests), ("max_items", max_items), ("max_errors", max_errors), ): _require_optional_positive(value, name) if max_depth is not None and max_depth < 0: msg = "max_depth must be non-negative" raise ValueError(msg) if min_items is not None and min_items < 0: msg = "min_items must be non-negative" raise ValueError(msg) _require_optional_rate(max_error_rate, "max_error_rate") _require_optional_rate(max_item_drop_rate, "max_item_drop_rate") max_duration_seconds = to_seconds(max_duration) if max_duration_seconds is not None and max_duration_seconds <= 0: msg = "max_duration must be positive" raise ValueError(msg) if metrics_port is not None and not 0 <= metrics_port <= 65535: msg = "metrics_port must be between 0 and 65535" raise ValueError(msg) self.spider = spider client: FetchClient = ( http_client if http_client is not None else HttpClient( concurrency=concurrency, emulation=emulation, timeout=request_timeout, html_max_size_bytes=html_max_size_bytes, keep_alive=keep_alive, max_response_size_bytes=max_response_size_bytes, ) ) require_positive_int(client.concurrency, "http_client.concurrency") self.http: FetchClient = ( http_cache.wrap(client) if http_cache is not None else client ) # Bound the queue to avoid unbounded growth when many requests are scheduled. default_queue_size = self.http.concurrency * 10 if max_pending_requests is not None: require_positive_int(max_pending_requests, "max_pending_requests") self.max_pending_requests: int = ( max_pending_requests if max_pending_requests is not None else default_queue_size ) self._request_order = count() # The queue itself is unbounded: ``_wait_for_queue_capacity`` enforces # ``max_pending_requests`` so it can let a worker overflow the bound # rather than deadlock (see that method). self._queue: asyncio.PriorityQueue[PrioritizedRequest] = asyncio.PriorityQueue() self._capacity_waiters: deque[tuple[asyncio.Future[bool], int | None]] = deque() self._stalled_workers: Counter[int] = Counter() self._worker_count = 0 # 16-byte digests of dedup keys; far smaller than the keys themselves. self._seen: set[bytes] = set() self.dedup_key: DedupKey = dedup_key or default_dedup_key self._stop_event = asyncio.Event() self.logger: Logger = get_logger(component="engine", spider=self.spider.name) self.engine_logger: EngineLogger = engine_logger or EngineLogger() self.request_middlewares: list[RequestMiddleware] = list( request_middlewares or [] ) self.response_middlewares: list[ResponseMiddleware] = list( response_middlewares or [] ) self.item_pipelines: list[ItemPipeline] = list(item_pipelines or []) self.item_batch_size = item_batch_size self.item_batch_wait: float = float(item_batch_wait) self._item_queue: ( asyncio.Queue[tuple[JSONValue, asyncio.Future[None]]] | None ) = None self._item_batch_ready = asyncio.Event() self._item_worker_task: asyncio.Task[None] | None = None self._lifecycle_closers: list[LifecycleCloser] = [] # Scheduling policy self._allowed_domains = normalize_domains( getattr(spider, "allowed_domains", ()) ) self._offsite_hosts_logged: set[str] = set() self._domain_slots = ( DomainSlots(concurrency_per_domain) if concurrency_per_domain is not None else None ) self.max_depth: int | None = max_depth self.max_requests: int | None = max_requests self.max_items: int | None = max_items self.max_errors: int | None = max_errors self.max_duration: float | None = max_duration_seconds self.max_error_rate: float | None = max_error_rate self.min_items: int | None = min_items self.max_item_drop_rate: float | None = max_item_drop_rate self._close_reason: str | None = None self._fetches_started = 0 self._items_reserved = 0 self._in_flight = 0 # Persistence and observability self._job_dir = job_dir self._job: JobState | None = None self._metrics_port = metrics_port self._metrics_host = metrics_host self.metrics_server: MetricsServer | None = None # Statistics tracking self.log_stats_interval: float | None = log_stats_interval self._start_time: float = 0.0 self._event_loop_type: str | None = None self.stats: CrawlStats = CrawlStats() self._stats: dict[str, int] = self.stats.counters
@property def close_reason(self) -> str | None: """Return why the crawl is stopping, or ``None`` while it runs normally.""" return self._close_reason @property def in_flight(self) -> int: """Return the number of requests currently being processed.""" return self._in_flight
[docs] def stop(self, reason: str = "shutdown") -> None: """Stop the crawl gracefully. New requests are no longer scheduled and queued ones are discarded (they stay saved when a job directory is configured, so the crawl can resume). Requests already being processed finish, pipelines close normally, and :meth:`run` returns with ``reason`` as the close reason. Calling it again has no effect. """ if self._close_reason is not None: return self._close_reason = reason self.logger.info( "Stopping crawl", spider=self.spider.name, reason=reason, pending_requests=self._queue.qsize(), in_flight=self._in_flight, ) # Let producers blocked on queue capacity observe the close. for waiter, _ in self._capacity_waiters: if not waiter.done(): waiter.set_result(True) self._capacity_waiters.clear() while True: try: self._queue.get_nowait() except asyncio.QueueEmpty: break self.stats.inc("dropped_requests") self._queue.task_done()
[docs] async def open_spider(self) -> None: """Open middleware, spider, and pipelines, then enqueue initial requests. Middleware opens before the spider; pipelines open afterward in their configured order. When resuming a job, saved requests are queued before ``start_requests()`` runs (already-seen start requests are skipped). """ if self._lifecycle_closers: raise RuntimeError("Spider lifecycle is already open") self.logger.info("Opening spider", spider=self.spider.name) try: await self._open_middlewares() self._register_lifecycle_close( f"spider {self.spider.name}", self.spider.close, ) await self.spider.open() for pipe in self.item_pipelines: self._register_lifecycle_close( f"pipeline {pipe.__class__.__name__}", lambda pipe=pipe: pipe.close(self.spider), ) await pipe.open(self.spider) if self.item_batch_size > 1: self._start_item_worker() self._register_lifecycle_close( "item batch coordinator", self._close_item_worker, ) self._restore_pending_requests() try: await self._run_callback( self.spider.start_requests, name="start_requests", url=None, response=None, parent=None, ) except _CrawlClosed: self.logger.debug( "Stopped start_requests because the crawl is closing", reason=self._close_reason, ) except BaseException as exc: cleanup_errors = await self._close_lifecycle_components() self._record_cleanup_failures(exc, cleanup_errors) raise
def _restore_pending_requests(self) -> None: if self._job is None or not self._job.resumed: return restored = self._job.load_pending() for seq, request in restored: self._queue.put_nowait((-request.priority, seq, request)) if restored: self._request_order = count(max(seq for seq, _ in restored) + 1) self.logger.info( "Resuming job", spider=self.spider.name, job_dir=str(self._job.directory), restored_requests=len(restored), seen_requests=self._job.seen_count(), )
[docs] async def close_spider(self) -> None: """Close pipelines, the spider, and middleware lifecycle hooks. Components close in reverse startup order. Middleware instances close once even when registered for both request and response work. """ self.logger.info("Closing spider", spider=self.spider.name) self._raise_cleanup_errors(await self._close_lifecycle_components())
def _iter_middlewares(self) -> Iterable[object]: seen_ids: set[int] = set() for middleware in [*self.request_middlewares, *self.response_middlewares]: middleware_id = id(middleware) if middleware_id in seen_ids: continue seen_ids.add(middleware_id) yield middleware def _iter_exception_middlewares(self) -> Iterable[ExceptionMiddleware]: for middleware in self._iter_middlewares(): process_exception = getattr(middleware, "process_exception", None) if callable(process_exception): yield cast("ExceptionMiddleware", middleware) async def _open_middlewares(self) -> None: for middleware in self._iter_middlewares(): close_hook = getattr(middleware, "close", None) if callable(close_hook): self._register_lifecycle_close( f"middleware {middleware.__class__.__name__}", lambda close_hook=close_hook: cast( "Awaitable[object]", close_hook(self.spider) ), ) open_hook = getattr(middleware, "open", None) if callable(open_hook): await cast("Awaitable[object]", open_hook(self.spider)) def _register_lifecycle_close( self, name: str, closer: Callable[[], Awaitable[object]], ) -> None: self._lifecycle_closers.append((name, closer)) async def _close_lifecycle_components(self) -> list[BaseException]: closers = self._lifecycle_closers self._lifecycle_closers = [] errors: list[BaseException] = [] for name, closer in reversed(closers): try: await closer() except BaseException as exc: errors.append(exc) self.logger.exception( "Lifecycle cleanup failed", component=name, error=str(exc), error_type=exc.__class__.__name__, ) return errors def _record_cleanup_failures( self, primary: BaseException, errors: Iterable[BaseException], ) -> None: for error in errors: primary.add_note(f"Cleanup failed with {error.__class__.__name__}: {error}") def _raise_cleanup_errors(self, errors: list[BaseException]) -> None: if not errors: return if len(errors) == 1: raise errors[0] raise BaseExceptionGroup("Multiple resource cleanup failures", errors) async def _shutdown(self, *, finished: bool) -> list[BaseException]: errors = await self._close_lifecycle_components() try: await self.http.close() except BaseException as exc: errors.append(exc) self.logger.exception( "HTTP client cleanup failed", error=str(exc), error_type=exc.__class__.__name__, ) if self._job is not None: job, self._job = self._job, None try: job.close(finished=finished) self.logger.info( "Saved job state", job_dir=str(job.directory), status="finished" if finished else "paused", ) except BaseException as exc: errors.append(exc) self.logger.exception( "Job state cleanup failed", error=str(exc), error_type=exc.__class__.__name__, ) if self.metrics_server is not None: try: await self.metrics_server.close() except BaseException as exc: # noqa: BLE001 - attempt every cleanup errors.append(exc) try: complete_logs() except BaseException as exc: # noqa: BLE001 - logging flush is final cleanup errors.append(exc) return errors async def _apply_request_mw(self, req: Request) -> Request: for mw in self.request_middlewares: req = await mw.process_request(req, self.spider) return req async def _handle_request_exception( self, req: Request, exc: Exception, ) -> bool: for mw in self._iter_exception_middlewares(): self.logger.debug( "Calling exception middleware", url=req.url, middleware=mw.__class__.__name__, error=str(exc), error_type=exc.__class__.__name__, ) retry_request = await mw.process_exception(req, exc, self.spider) if retry_request is None: self.logger.debug( "Exception middleware did not retry request", url=req.url, middleware=mw.__class__.__name__, error_type=exc.__class__.__name__, ) continue if not isinstance(retry_request, Request): self.logger.warning( "Ignoring invalid exception middleware result", url=req.url, middleware=mw.__class__.__name__, result_type=retry_request.__class__.__name__, error_type=exc.__class__.__name__, ) continue self.engine_logger.retrying_request( self.logger, retry_request, self.spider, source=f"exception middleware {mw.__class__.__name__}", ) self.stats.inc("retries") await self._enqueue(retry_request) return True return False async def _handle_request_errback( self, req: Request, exc: Exception, ) -> bool: errback = req.errback if errback is None: return False name = getattr(errback, "__name__", errback.__class__.__name__) self.logger.debug( "Calling request errback", url=req.url, errback=name, error=str(exc), error_type=exc.__class__.__name__, ) await self._run_callback( lambda: errback(req, exc), name=name, url=req.url, response=None, parent=req, ) return True async def _enqueue(self, req: Request, parent: Request | None = None) -> None: """Filter, deduplicate, and queue ``req`` (scheduled from ``parent``).""" req = self._with_depth(req, parent) if not self._passes_scheduling_filters(req): return # Serialize before marking the request seen so an unrestorable request # fails loudly without polluting the seen-set. payload = self._job.serialize(req) if self._job is not None else None if not req.dont_filter and not self._mark_seen(req): self.stats.inc("dupe_filtered") self.logger.debug("Skipping already seen request", url=req.url) return if self._close_reason is None: await self._wait_for_queue_capacity(req) seq = next(self._request_order) if self._job is not None and payload is not None: self._job.add_pending(seq, req, payload) if self._close_reason is not None: self.stats.inc("dropped_requests") if _CURRENT_WORKER.get() is None: raise _CrawlClosed return self._queue.put_nowait((-req.priority, seq, req)) self.logger.debug( "Enqueued request", url=req.url, dont_filter=req.dont_filter, priority=req.priority, depth=_meta_depth(req), ) def _with_depth(self, req: Request, parent: Request | None) -> Request: if parent is not None: depth = _meta_depth(parent) + 1 elif "depth" in req.meta: return req else: depth = 0 return req.replace(meta={**req.meta, "depth": depth}) def _passes_scheduling_filters(self, req: Request) -> bool: if self._allowed_domains and not req.dont_filter: host = url_host(req.url) if not host_in_domains(host, self._allowed_domains): self.stats.inc("offsite_filtered") if ( host not in self._offsite_hosts_logged and len(self._offsite_hosts_logged) < 1000 ): self._offsite_hosts_logged.add(host) self.logger.debug( "Filtered offsite request", url=req.url, host=host, allowed_domains=list(self._allowed_domains), ) return False if self.max_depth is not None and _meta_depth(req) > self.max_depth: self.stats.inc("depth_filtered") self.logger.debug( "Filtered request beyond max_depth", url=req.url, depth=_meta_depth(req), max_depth=self.max_depth, ) return False return True def _mark_seen(self, req: Request) -> bool: """Record ``req``'s dedup key; return ``False`` when already seen.""" digest = hashlib.blake2b(self.dedup_key(req).encode(), digest_size=16).digest() if self._job is not None: return self._job.seen_add(digest) if digest in self._seen: return False self._seen.add(digest) return True def _seen_count(self) -> int: return self._job.seen_count() if self._job is not None else len(self._seen) async def _wait_for_queue_capacity(self, req: Request) -> None: """Apply ``max_pending_requests`` backpressure to one enqueue. Requests scheduled outside workers (``start_requests()``) wait until the queue has room. Workers are also the queue's only consumers, so a worker's callback (or a task it spawned) may wait only while no other worker is waiting and at least one other worker exists. That throttles a single heavy producer while the remaining workers keep consuming. Any other worker that finds the queue full enqueues past the bound and releases waiting workers to do the same. When the frontier grows faster than it is consumed, blocking workers would only idle them (and keep their responses alive) without bounding the queue, so they keep crawling instead. This also guarantees the workers never deadlock. """ worker = _CURRENT_WORKER.get() while self._queue.qsize() >= self.max_pending_requests: if worker is not None and not self._worker_may_wait(worker): self.logger.debug( "Queue full; enqueuing past max_pending_requests to keep " "workers crawling", url=req.url, queue_size=self._queue.qsize(), max_pending_requests=self.max_pending_requests, ) self._release_waiting_workers() return waiter: asyncio.Future[bool] = asyncio.get_running_loop().create_future() self._capacity_waiters.append((waiter, worker)) if worker is not None: self._stalled_workers[worker] += 1 try: overflow = await waiter except asyncio.CancelledError: if waiter.done() and not waiter.cancelled() and not waiter.result(): # Woken for a free slot but cancelled first; pass it on. self._wake_capacity_waiter() raise finally: with suppress(ValueError): self._capacity_waiters.remove((waiter, worker)) if worker is not None: self._stalled_workers[worker] -= 1 if not self._stalled_workers[worker]: del self._stalled_workers[worker] if overflow: return def _worker_may_wait(self, worker: int) -> bool: if self._worker_count < 2: return False return all(stalled == worker for stalled in self._stalled_workers) def _wake_capacity_waiter(self) -> None: """Wake the oldest waiter because a queue slot was freed.""" while self._capacity_waiters: waiter, _ = self._capacity_waiters.popleft() if not waiter.done(): waiter.set_result(False) return def _release_waiting_workers(self) -> None: """Let every waiting worker enqueue past the bound.""" remaining: deque[tuple[asyncio.Future[bool], int | None]] = deque() for waiter, worker in self._capacity_waiters: if worker is not None and not waiter.done(): waiter.set_result(True) elif not waiter.done(): remaining.append((waiter, worker)) self._capacity_waiters = remaining async def _worker(self, index: int = 0) -> None: # Each worker runs in its own task, so this only tags this worker and the # tasks its callbacks spawn. _CURRENT_WORKER.set(index) while not self._stop_event.is_set(): try: async with asyncio.timeout(1.0): _, seq, req = await self._queue.get() except TimeoutError: if self._stop_event.is_set(): break continue except asyncio.CancelledError: break self._wake_capacity_waiter() completed = True try: if self._close_reason is not None or not self._reserve_fetch(): # Stopping: leave the request journaled so a job can resume. completed = False self.stats.inc("dropped_requests") continue self._in_flight += 1 try: await self._process_request(req) finally: self._in_flight -= 1 except IgnoreRequest as exc: self.stats.inc("ignored_requests") self.stats.inc_labeled("ignored_by_reason", exc.reason) self.logger.debug( "Ignored request", url=req.url, reason=exc.reason, detail=str(exc), ) except CloseSpider as exc: self.stop(exc.reason) except Exception as exc: # noqa: BLE001 - logged by the failure handler await self._handle_request_failure(req, exc) finally: if self._job is not None and completed: self._job.remove_pending(seq) self._queue.task_done() def _reserve_fetch(self) -> bool: """Claim one of ``max_requests`` fetches; stop the crawl when exhausted.""" if self.max_requests is None: return True if self._fetches_started >= self.max_requests: self.stop("max_requests") return False self._fetches_started += 1 return True async def _process_request(self, req: Request) -> None: req = await self._apply_request_mw(req) self.engine_logger.fetching_request(self.logger, req, self.spider) self.stats.inc("requests_sent") self.stats.inc_labeled("requests_by_domain", url_host(req.url) or "-") if self._domain_slots is not None: async with self._domain_slots.acquire(req.url): resp = await self.http.fetch(req) else: resp = await self.http.fetch(req) self.stats.inc("responses_received") self.stats.inc_labeled("responses_by_status", str(resp.status)) self.engine_logger.fetched_response(self.logger, req, resp, self.spider) await self._handle_response(resp) async def _handle_request_failure(self, req: Request, exc: Exception) -> None: if await self._handle_request_exception(req, exc): return self.stats.inc("errors") cause = exc.__cause__ if isinstance(exc, SpiderError) else None self.stats.inc_labeled("errors_by_type", type(cause or exc).__name__) if self.max_errors is not None and self.stats.get("errors") >= self.max_errors: self.stop("max_errors") try: if await self._handle_request_errback(req, exc): return except Exception as errback_exc: errback_cause = errback_exc.__cause__ or errback_exc.__context__ error_context = { "url": req.url, "error": str(errback_exc), "error_type": errback_exc.__class__.__name__, "original_error": str(exc), "original_error_type": exc.__class__.__name__, "spider": self.spider.name, } if errback_cause is not None: error_context["cause"] = self._safe_repr(errback_cause) error_context["cause_type"] = errback_cause.__class__.__name__ self.logger.error( "Request errback failed", **error_context, exc_info=not isinstance(errback_exc, SilkwormError), ) return failure_cause = exc.__cause__ or exc.__context__ error_context = { "url": req.url, "error": str(exc), "error_type": exc.__class__.__name__, "spider": self.spider.name, } if failure_cause is not None: error_context["cause"] = self._safe_repr(failure_cause) error_context["cause_type"] = failure_cause.__class__.__name__ # silkworm's own errors are self-explanatory (and callback failures are # already logged with a traceback); only unexpected errors, e.g. bugs in # middlewares or pipelines, get one here. self.logger.error( "Failed to process request", **error_context, exc_info=not isinstance(exc, SilkwormError), ) async def _apply_response_mw( self, resp: Response, owned_responses: list[Response], ) -> Response | Request: current: Response | Request = resp for mw in self.response_middlewares: if isinstance(current, Request): # already converted to a retry Request by a previous mw break previous = current current = await mw.process_response(previous, self.spider) if isinstance(current, Response) and current is not previous: owned_responses.append(current) return current async def _handle_response(self, resp: Response) -> None: owned_responses = [resp] try: processed = await self._apply_response_mw(resp, owned_responses) if isinstance(processed, Request): # e.g. RetryMiddleware wants a retry self.engine_logger.retrying_request( self.logger, processed, self.spider, source="response middleware", ) self.stats.inc("retries") await self._enqueue(processed) return callback = processed.request.callback name = getattr(callback, "__name__", "parse") if callback else "parse" effective_callback = callback or self.spider.parse if self._expects_html(callback): callback_resp = self._ensure_html_response(processed) if callback_resp is not processed: owned_responses.append(callback_resp) else: callback_resp = processed await self._run_callback( lambda: effective_callback(callback_resp), name=name, url=processed.url, response=callback_resp, parent=processed.request, ) finally: primary = sys.exception() closed_ids: set[int] = set() cleanup_errors: list[BaseException] = [] for owned_response in reversed(owned_responses): if id(owned_response) in closed_ids: continue closed_ids.add(id(owned_response)) try: owned_response.close() except BaseException as exc: # noqa: BLE001 - close every response cleanup_errors.append(exc) if primary is not None: self._record_cleanup_failures(primary, cleanup_errors) else: self._raise_cleanup_errors(cleanup_errors) async def _run_callback( self, invoke: Callable[[], object], *, name: str, url: str | None, response: Response | None, parent: Request | None, ) -> None: """Run a callback coroutine inside a scope wired to the engine sinks. Items reported with ``emit`` reach the pipelines and requests reported with ``follow`` reach the queue (one level deeper than ``parent``) while the callback runs. :class:`~silkworm.exceptions.CloseSpider` stops the crawl instead of failing the callback. """ scope = CrawlScope( owner=name, emit_item=self._emit_item, schedule_request=partial(self._enqueue, parent=parent), flush_items=self._request_item_flush, response=response, ) try: async with enter_scope(scope): await invoke_callback(invoke, name) except CallbackContractError: raise except CloseSpider as exc: self.stop(exc.reason) except Exception as exc: raise self._callback_failure(name, url, exc) from exc def _callback_failure( self, name: str, url: str | None, exc: Exception, ) -> SpiderError: self.logger.exception( "Spider callback failed", callback=name, spider=self.spider.name, url=url, error=str(exc), error_type=exc.__class__.__name__, ) return SpiderError(f"Spider callback '{name}' failed for {self.spider.name}") async def _emit_item(self, item: JSONLike) -> Awaitable[None] | None: if self.max_items is not None and self._items_reserved >= self.max_items: self._count_dropped_item(MAX_ITEMS_DROP_REASON) return None self._items_reserved += 1 if self._item_queue is not None: completion: asyncio.Future[None] = ( asyncio.get_running_loop().create_future() ) await self._item_queue.put((cast(JSONValue, item), completion)) self._item_batch_ready.set() return completion accepted = False try: self.logger.debug( "Processing scraped item", spider=self.spider.name, pipelines=len(self.item_pipelines), ) # Pipelines take JSONValue; callbacks may emit read-only JSONLike # shapes, which are the same objects at runtime. accepted = bool(await self._process_item(cast(JSONValue, item))) finally: if not accepted: self._items_reserved -= 1 if ( accepted and self.max_items is not None and self.stats.get("items_scraped") >= self.max_items ): self.stop(MAX_ITEMS_DROP_REASON) return None def _start_item_worker(self) -> None: self._item_queue = asyncio.Queue(maxsize=self.item_batch_size * 10) self._item_worker_task = asyncio.create_task(self._item_batch_worker()) def _request_item_flush(self) -> None: if self._item_queue is not None: self._item_batch_ready.set() async def _close_item_worker(self) -> None: queue = self._item_queue task = self._item_worker_task if queue is None or task is None: return self._item_batch_ready.set() await queue.join() task.cancel() with suppress(asyncio.CancelledError): await task self._item_queue = None self._item_worker_task = None async def _item_batch_worker(self) -> None: queue = self._item_queue assert queue is not None while True: first = await queue.get() batch = [first] deadline = asyncio.get_running_loop().time() + self.item_batch_wait while len(batch) < self.item_batch_size: while len(batch) < self.item_batch_size: try: batch.append(queue.get_nowait()) except asyncio.QueueEmpty: break if len(batch) >= self.item_batch_size: break self._item_batch_ready.clear() if not queue.empty(): continue remaining = deadline - asyncio.get_running_loop().time() if remaining <= 0: break try: async with asyncio.timeout(remaining): await self._item_batch_ready.wait() except TimeoutError: break items = [item for item, _ in batch] try: accepted = await self._process_items(items) except Exception as exc: # noqa: BLE001 - report one failure to every waiter self._items_reserved -= len(batch) for _, completion in batch: if not completion.done(): completion.set_exception(exc) else: self._items_reserved -= len(batch) - accepted for _, completion in batch: if not completion.done(): completion.set_result(None) if ( self.max_items is not None and self.stats.get("items_scraped") >= self.max_items ): self.stop(MAX_ITEMS_DROP_REASON) finally: for _ in batch: queue.task_done() async def _process_item(self, item: JSONValue) -> bool: """Run ``item`` through the pipelines; return whether it was kept.""" for pipe in self.item_pipelines: self.engine_logger.running_item_pipeline( self.logger, pipe, self.spider, ) try: item = await pipe.process_item(item, self.spider) except DropItem as exc: self._count_dropped_item(exc.reason) self.logger.debug( "Dropped item", pipeline=pipe.__class__.__name__, reason=exc.reason, detail=str(exc), ) return False self.stats.inc("items_scraped") return True async def _process_items(self, items: list[JSONValue]) -> int: """Run a batch through pipelines, using native bulk paths when safe.""" active = items for pipe in self.item_pipelines: self.engine_logger.running_item_pipeline( self.logger, pipe, self.spider, ) if getattr(pipe, "native_batch", False): processed = await pipe.process_items(active, self.spider) # type: ignore[attr-defined] if len(processed) != len(active): raise RuntimeError( f"{pipe.__class__.__name__}.process_items() must return " "one item for every input when native_batch is enabled" ) active = processed continue processed = [] for item in active: try: processed.append(await pipe.process_item(item, self.spider)) except DropItem as exc: self._count_dropped_item(exc.reason) self.logger.debug( "Dropped item", pipeline=pipe.__class__.__name__, reason=exc.reason, detail=str(exc), ) active = processed self.stats.inc("items_scraped", len(active)) return len(active) def _count_dropped_item(self, reason: str) -> None: self.stats.inc("items_dropped") self.stats.inc_labeled("items_dropped_by_reason", reason) def _get_memory_usage_mb(self) -> float: """ Return memory usage (RSS) in megabytes, normalizing platform differences. """ if resource is None: return 0.0 usage = resource.getrusage(resource.RUSAGE_SELF).ru_maxrss divisor = 1024 * 1024 if sys.platform == "darwin" else 1024 return usage / divisor def _detect_event_loop(self, loop: asyncio.AbstractEventLoop | None = None) -> str: """ Identify which event loop implementation is currently running. """ loop = loop or asyncio.get_running_loop() module = loop.__class__.__module__.lower() name = loop.__class__.__name__.lower() if "rsloop" in module or "rsloop" in name: return "rsloop" if "uvloop" in module or "uvloop" in name: return "uvloop" if "trio" in module or "trio" in name: return "trio" return "asyncio" def _stats_payload(self, elapsed: float) -> dict[str, object]: requests_rate = self.stats.get("requests_sent") / elapsed if elapsed > 0 else 0 payload: dict[str, object] = { "elapsed_seconds": round(elapsed, 1), **self.stats.counters, "queue_size": self._queue.qsize(), "in_flight": self._in_flight, "requests_per_second": round(requests_rate, 2), "seen_requests": self._seen_count(), "memory_mb": round(self._get_memory_usage_mb(), 2), } for name, counter in self.stats.labeled.items(): if counter: payload[name] = dict(counter) return payload def _statistics_log_context( self, elapsed: float, *, include_event_loop: bool = False, ) -> dict[str, object]: context: dict[str, object] = { "spider": self.spider.name, **self._stats_payload(elapsed), **self.spider.stats_payload, } if include_event_loop: context["event_loop"] = self._event_loop_type return context
[docs] def metrics_text(self) -> str: """Return current statistics in the Prometheus text exposition format.""" elapsed = time.time() - self._start_time if self._start_time else 0.0 return render_metrics( spider=self.spider.name, stats=self.stats, gauges={ "queue_size": self._queue.qsize(), "in_flight": self._in_flight, "seen_requests": self._seen_count(), "elapsed_seconds": round(elapsed, 3), "memory_mb": round(self._get_memory_usage_mb(), 2), "running": 0 if self._stop_event.is_set() else 1, }, custom=self.spider.stats_payload, )
async def _log_statistics(self) -> None: """Periodically log statistics about the crawl progress.""" if self.log_stats_interval is None: return interval = self.log_stats_interval if interval <= 0: return while not self._stop_event.is_set(): try: async with asyncio.timeout(interval): await self._stop_event.wait() break except TimeoutError: self.logger.info( "Crawl statistics", **self._statistics_log_context(time.time() - self._start_time), ) async def _enforce_max_duration(self, seconds: float) -> None: try: async with asyncio.timeout(seconds): await self._stop_event.wait() except TimeoutError: self.stop("max_duration") def _evaluate_failure_policy(self, close_reason: str) -> tuple[str, ...]: if close_reason == "shutdown": return () failures: list[str] = [] sent = self.stats.get("requests_sent") errors = self.stats.get("errors") if ( self.max_error_rate is not None and sent and errors / sent > self.max_error_rate ): failures.append( f"error rate {errors / sent:.1%} exceeds max_error_rate " f"{self.max_error_rate:.1%} ({errors} errors / {sent} requests)" ) scraped = self.stats.get("items_scraped") if self.min_items is not None and scraped < self.min_items: failures.append( f"scraped {scraped} items, fewer than min_items={self.min_items}" ) if self.max_item_drop_rate is not None: dropped = self.stats.get("items_dropped") - self.stats.labeled[ "items_dropped_by_reason" ].get(MAX_ITEMS_DROP_REASON, 0) total = scraped + dropped if total and dropped / total > self.max_item_drop_rate: failures.append( f"item drop rate {dropped / total:.1%} exceeds " f"max_item_drop_rate {self.max_item_drop_rate:.1%} " f"({dropped} dropped / {total} items)" ) return tuple(failures) def _final_statistics(self, close_reason: str | None) -> None: self.logger.info( "Final crawl statistics", close_reason=close_reason, **self._statistics_log_context( time.time() - self._start_time, include_event_loop=True, ), )
[docs] async def run(self) -> CrawlResult: """Run the crawl until the queue drains or a stop condition, then clean up. Worker tasks, periodic statistics, and lifecycle hooks are managed as a task group. The HTTP client, spider components, job state, and metrics server are always closed, and a final statistics record is emitted. Returns: The crawl's :class:`~silkworm.CrawlResult`. Raises: CrawlFailedError: If the crawl violated its failure policy (``max_error_rate``, ``min_items``, ``max_item_drop_rate``). Exception: Errors from lifecycle hooks, pipelines' ``open``/``close``, or cleanup propagate after structured error logging. """ self.logger.info("Starting engine", spider=self.spider.name) self._start_time = time.time() self._event_loop_type = self._detect_event_loop() try: if self._job_dir is not None and self._job is None: self._job = JobState(self._job_dir, self.spider) if self._metrics_port is not None: self.metrics_server = MetricsServer( self.metrics_text, host=self._metrics_host, port=self._metrics_port, ) await self.metrics_server.start() async with asyncio.TaskGroup() as tg: self._worker_count = self.http.concurrency for index in range(self._worker_count): tg.create_task(self._worker(index)) if self.log_stats_interval is not None and self.log_stats_interval > 0: tg.create_task(self._log_statistics()) if self.max_duration is not None: tg.create_task(self._enforce_max_duration(self.max_duration)) # Open spider and seed initial requests while workers are already waiting. await self.open_spider() await self._queue.join() self._stop_event.set() except BaseException as exc: self._stop_event.set() self._final_statistics(self._close_reason or "error") cleanup_errors = await self._shutdown(finished=False) self._record_cleanup_failures(exc, cleanup_errors) raise close_reason = self._close_reason or "finished" self._stop_event.set() self._final_statistics(close_reason) self._raise_cleanup_errors( await self._shutdown(finished=close_reason == "finished") ) result = freeze_result( spider=self.spider.name, close_reason=close_reason, elapsed_seconds=time.time() - self._start_time, stats=self.stats, custom_stats=self.spider.stats_payload, failures=self._evaluate_failure_policy(close_reason), ) if result.failures: self.logger.error( "Crawl failed its failure policy", spider=self.spider.name, failures=list(result.failures), ) raise CrawlFailedError(result) return result
def _expects_html(self, callback: Callback | None) -> bool: if callback is None: return True cb_self = getattr(callback, "__self__", None) cb_func = getattr(callback, "__func__", None) parse_func = getattr(self.spider.parse, "__func__", None) return cb_self is self.spider and cb_func is parse_func def _ensure_html_response(self, resp: Response) -> HTMLResponse: if isinstance(resp, HTMLResponse): return resp return HTMLResponse( url=resp.url, status=resp.status, headers=resp.headers, body=resp.body, request=resp.request, doc_max_size_bytes=self.http.html_max_size_bytes, ) def _safe_repr(self, value: object, limit: int = 200) -> str: if value is None: return "None" text = _SAFE_REPR.repr(value) return text if len(text) <= limit else f"{text[:limit]}..."