Skip to content

Commit f518d97

Browse files
committed
Move utils
1 parent 82f0ed9 commit f518d97

6 files changed

Lines changed: 251 additions & 232 deletions

File tree

pylintrc

Lines changed: 0 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -10,7 +10,6 @@ disable=
1010
too-few-public-methods,
1111
too-many-arguments,
1212
too-many-instance-attributes,
13-
too-many-lines,
1413
# tests
1514
duplicate-code,
1615
import-outside-toplevel,

scrapy_playwright/_loop.py

Lines changed: 107 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,107 @@
1+
import asyncio
2+
import platform
3+
from dataclasses import dataclass
4+
from threading import Thread
5+
from typing import Awaitable, Dict
6+
7+
from twisted.internet.defer import Deferred
8+
from twisted.python import failure
9+
10+
from scrapy_playwright._utils import logger
11+
12+
13+
@dataclass
14+
class _QueueItem:
15+
coro: Awaitable
16+
promise: Deferred | asyncio.Future
17+
loop: asyncio.AbstractEventLoop | None = None
18+
19+
20+
class _ThreadedLoopAdapter:
21+
"""Utility class to start an asyncio event loop in a new thread and redirect coroutines.
22+
This allows to run Playwright in a different loop than the Scrapy crawler, allowing to
23+
use ProactorEventLoop which is supported by Playwright on Windows.
24+
"""
25+
26+
_loop: asyncio.AbstractEventLoop
27+
_thread: Thread
28+
_coro_queue: asyncio.Queue = asyncio.Queue()
29+
_stop_events: Dict[int, asyncio.Event] = {}
30+
31+
@classmethod
32+
async def _handle_coro_deferred(cls, queue_item: _QueueItem) -> None:
33+
from twisted.internet import reactor
34+
35+
dfd: Deferred = queue_item.promise
36+
37+
try:
38+
result = await queue_item.coro
39+
except Exception as exc:
40+
reactor.callFromThread(dfd.errback, failure.Failure(exc))
41+
else:
42+
reactor.callFromThread(dfd.callback, result)
43+
44+
@classmethod
45+
async def _handle_coro_future(cls, queue_item: _QueueItem) -> None:
46+
future: asyncio.Future = queue_item.promise
47+
loop: asyncio.AbstractEventLoop = queue_item.loop # type: ignore[assignment]
48+
try:
49+
result = await queue_item.coro
50+
except Exception as exc:
51+
loop.call_soon_threadsafe(future.set_exception, exc)
52+
else:
53+
loop.call_soon_threadsafe(future.set_result, result)
54+
55+
@classmethod
56+
async def _process_queue(cls) -> None:
57+
while any(not ev.is_set() for ev in cls._stop_events.values()):
58+
queue_item = await cls._coro_queue.get()
59+
if isinstance(queue_item.promise, asyncio.Future):
60+
asyncio.create_task(cls._handle_coro_future(queue_item))
61+
elif isinstance(queue_item.promise, Deferred):
62+
asyncio.create_task(cls._handle_coro_deferred(queue_item))
63+
cls._coro_queue.task_done()
64+
65+
@classmethod
66+
def _deferred_from_coro(cls, coro: Awaitable) -> Deferred:
67+
dfd: Deferred = Deferred()
68+
queue_item = _QueueItem(coro=coro, promise=dfd)
69+
asyncio.run_coroutine_threadsafe(cls._coro_queue.put(queue_item), cls._loop)
70+
return dfd
71+
72+
@classmethod
73+
def _future_from_coro(cls, coro: Awaitable) -> asyncio.Future:
74+
target_loop = asyncio.get_running_loop() # Scrapy thread loop
75+
future: asyncio.Future = asyncio.Future()
76+
queue_item = _QueueItem(coro=coro, promise=future, loop=target_loop)
77+
asyncio.run_coroutine_threadsafe(cls._coro_queue.put(queue_item), cls._loop)
78+
return future
79+
80+
@classmethod
81+
def start(cls, download_handler_id: int) -> None:
82+
"""Start the event loop in a new thread if not already started.
83+
Should be called from the Scrapy thread.
84+
"""
85+
cls._stop_events[download_handler_id] = asyncio.Event()
86+
if not getattr(cls, "_loop", None):
87+
policy = asyncio.DefaultEventLoopPolicy()
88+
if platform.system() == "Windows":
89+
policy = asyncio.WindowsProactorEventLoopPolicy() # type: ignore[attr-defined]
90+
cls._loop = policy.new_event_loop()
91+
92+
if not getattr(cls, "_thread", None):
93+
cls._thread = Thread(target=cls._loop.run_forever, daemon=True)
94+
cls._thread.start()
95+
logger.info("Started loop on separate thread: %s", cls._loop)
96+
asyncio.run_coroutine_threadsafe(cls._process_queue(), cls._loop)
97+
98+
@classmethod
99+
def stop(cls, download_handler_id: int) -> None:
100+
"""Wait until all handlers are closed to stop the event loop and join the thread.
101+
Should be called from the Scrapy thread.
102+
"""
103+
cls._stop_events[download_handler_id].set()
104+
if all(ev.is_set() for ev in cls._stop_events.values()):
105+
asyncio.run_coroutine_threadsafe(cls._coro_queue.join(), cls._loop)
106+
cls._loop.call_soon_threadsafe(cls._loop.stop)
107+
cls._thread.join()

scrapy_playwright/_utils.py

Lines changed: 129 additions & 102 deletions
Original file line numberDiff line numberDiff line change
@@ -1,17 +1,18 @@
1-
import asyncio
21
import logging
3-
import platform
4-
import threading
5-
from dataclasses import dataclass
6-
from typing import Awaitable, Dict, Iterator, Optional, Tuple, Union
2+
from typing import Awaitable, Callable, Iterator, Optional, Tuple, Union
73

8-
import scrapy
9-
from playwright.async_api import Error, Page, Request, Response
4+
from playwright.async_api import (
5+
Error,
6+
Page,
7+
Request as PlaywrightRequest,
8+
Response as PlaywrightResponse,
9+
)
10+
from scrapy import Spider
11+
from scrapy.http import Request as ScrapyRequest
1012
from scrapy.http.headers import Headers
1113
from scrapy.settings import Settings
14+
from scrapy.utils.misc import load_object
1215
from scrapy.utils.python import to_unicode
13-
from twisted.internet.defer import Deferred
14-
from twisted.python import failure
1516
from w3lib.encoding import html_body_declared_encoding, http_content_type_encoding
1617

1718

@@ -60,7 +61,7 @@ def _is_safe_close_error(error: Error) -> bool:
6061

6162
async def _get_page_content(
6263
page: Page,
63-
spider: scrapy.Spider,
64+
spider: Spider,
6465
context_name: str,
6566
scrapy_request_url: str,
6667
scrapy_request_method: str,
@@ -96,7 +97,7 @@ def _get_float_setting(settings: Settings, key: str) -> Optional[float]:
9697

9798

9899
async def _get_header_value(
99-
resource: Union[Request, Response],
100+
resource: Union[PlaywrightRequest, PlaywrightResponse],
100101
header_name: str,
101102
) -> Optional[str]:
102103
try:
@@ -105,98 +106,124 @@ async def _get_header_value(
105106
return None
106107

107108

108-
@dataclass
109-
class _QueueItem:
110-
coro: Awaitable
111-
promise: Deferred | asyncio.Future
112-
loop: asyncio.AbstractEventLoop | None = None
113-
114-
115-
class _ThreadedLoopAdapter:
116-
"""Utility class to start an asyncio event loop in a new thread and redirect coroutines.
117-
This allows to run Playwright in a different loop than the Scrapy crawler, allowing to
118-
use ProactorEventLoop which is supported by Playwright on Windows.
119-
"""
120-
121-
_loop: asyncio.AbstractEventLoop
122-
_thread: threading.Thread
123-
_coro_queue: asyncio.Queue = asyncio.Queue()
124-
_stop_events: Dict[int, asyncio.Event] = {}
125-
126-
@classmethod
127-
async def _handle_coro_deferred(cls, queue_item: _QueueItem) -> None:
128-
from twisted.internet import reactor
109+
def _attach_page_event_handlers(
110+
page: Page, request: ScrapyRequest, spider: Spider, context_name: str
111+
) -> None:
112+
event_handlers = request.meta.get("playwright_page_event_handlers") or {}
113+
for event, handler in event_handlers.items():
114+
if callable(handler):
115+
page.on(event, handler)
116+
elif isinstance(handler, str):
117+
try:
118+
page.on(event, getattr(spider, handler))
119+
except AttributeError as ex:
120+
logger.warning(
121+
"Spider '%s' does not have a '%s' attribute,"
122+
" ignoring handler for event '%s'",
123+
spider.name,
124+
handler,
125+
event,
126+
extra={
127+
"spider": spider,
128+
"context_name": context_name,
129+
"scrapy_request_url": request.url,
130+
"scrapy_request_method": request.method,
131+
"exception": ex,
132+
},
133+
exc_info=True,
134+
)
135+
136+
137+
async def _set_redirect_meta(request: ScrapyRequest, response: PlaywrightResponse) -> None:
138+
"""Update a Scrapy request with metadata about redirects."""
139+
redirect_times: int = 0
140+
redirect_urls: list = []
141+
redirect_reasons: list = []
142+
redirected = response.request.redirected_from
143+
while redirected is not None:
144+
redirect_times += 1
145+
redirect_urls.append(redirected.url)
146+
redirected_response = await redirected.response()
147+
reason = None if redirected_response is None else redirected_response.status
148+
redirect_reasons.append(reason)
149+
redirected = redirected.redirected_from
150+
if redirect_times:
151+
request.meta["redirect_times"] = redirect_times
152+
request.meta["redirect_urls"] = list(reversed(redirect_urls))
153+
request.meta["redirect_reasons"] = list(reversed(redirect_reasons))
154+
155+
156+
async def _maybe_execute_page_init_callback(
157+
page: Page,
158+
request: ScrapyRequest,
159+
context_name: str,
160+
spider: Spider,
161+
) -> None:
162+
page_init_callback = request.meta.get("playwright_page_init_callback")
163+
if page_init_callback:
164+
try:
165+
page_init_callback = load_object(page_init_callback)
166+
await page_init_callback(page, request)
167+
except Exception as ex:
168+
logger.warning(
169+
"[Context=%s] Page init callback exception for %s exc_type=%s exc_msg=%s",
170+
context_name,
171+
repr(request),
172+
type(ex),
173+
str(ex),
174+
extra={
175+
"spider": spider,
176+
"context_name": context_name,
177+
"scrapy_request_url": request.url,
178+
"scrapy_request_method": request.method,
179+
"exception": ex,
180+
},
181+
exc_info=True,
182+
)
129183

130-
dfd: Deferred = queue_item.promise
131184

132-
try:
133-
result = await queue_item.coro
134-
except Exception as exc:
135-
reactor.callFromThread(dfd.errback, failure.Failure(exc))
185+
def _make_request_logger(context_name: str, spider: Spider) -> Callable:
186+
async def _log_request(request: PlaywrightRequest) -> None:
187+
log_args = [context_name, request.method.upper(), request.url, request.resource_type]
188+
referrer = await _get_header_value(request, "referer")
189+
if referrer:
190+
log_args.append(referrer)
191+
log_msg = "[Context=%s] Request: <%s %s> (resource type: %s, referrer: %s)"
136192
else:
137-
reactor.callFromThread(dfd.callback, result)
138-
139-
@classmethod
140-
async def _handle_coro_future(cls, queue_item: _QueueItem) -> None:
141-
future: asyncio.Future = queue_item.promise
142-
loop: asyncio.AbstractEventLoop = queue_item.loop # type: ignore[assignment]
143-
try:
144-
result = await queue_item.coro
145-
except Exception as exc:
146-
loop.call_soon_threadsafe(future.set_exception, exc)
193+
log_msg = "[Context=%s] Request: <%s %s> (resource type: %s)"
194+
logger.debug(
195+
log_msg,
196+
*log_args,
197+
extra={
198+
"spider": spider,
199+
"context_name": context_name,
200+
"playwright_request_url": request.url,
201+
"playwright_request_method": request.method,
202+
"playwright_resource_type": request.resource_type,
203+
},
204+
)
205+
206+
return _log_request
207+
208+
209+
def _make_response_logger(context_name: str, spider: Spider) -> Callable:
210+
async def _log_response(response: PlaywrightResponse) -> None:
211+
log_args = [context_name, response.status, response.url]
212+
location = await _get_header_value(response, "location")
213+
if location:
214+
log_args.append(location)
215+
log_msg = "[Context=%s] Response: <%i %s> (location: %s)"
147216
else:
148-
loop.call_soon_threadsafe(future.set_result, result)
149-
150-
@classmethod
151-
async def _process_queue(cls) -> None:
152-
while any(not ev.is_set() for ev in cls._stop_events.values()):
153-
queue_item = await cls._coro_queue.get()
154-
if isinstance(queue_item.promise, asyncio.Future):
155-
asyncio.create_task(cls._handle_coro_future(queue_item))
156-
elif isinstance(queue_item.promise, Deferred):
157-
asyncio.create_task(cls._handle_coro_deferred(queue_item))
158-
cls._coro_queue.task_done()
159-
160-
@classmethod
161-
def _deferred_from_coro(cls, coro: Awaitable) -> Deferred:
162-
dfd: Deferred = Deferred()
163-
queue_item = _QueueItem(coro=coro, promise=dfd)
164-
asyncio.run_coroutine_threadsafe(cls._coro_queue.put(queue_item), cls._loop)
165-
return dfd
166-
167-
@classmethod
168-
def _future_from_coro(cls, coro: Awaitable) -> asyncio.Future:
169-
target_loop = asyncio.get_running_loop() # Scrapy thread loop
170-
future: asyncio.Future = asyncio.Future()
171-
queue_item = _QueueItem(coro=coro, promise=future, loop=target_loop)
172-
asyncio.run_coroutine_threadsafe(cls._coro_queue.put(queue_item), cls._loop)
173-
return future
174-
175-
@classmethod
176-
def start(cls, download_handler_id: int) -> None:
177-
"""Start the event loop in a new thread if not already started.
178-
Should be called from the Scrapy thread.
179-
"""
180-
cls._stop_events[download_handler_id] = asyncio.Event()
181-
if not getattr(cls, "_loop", None):
182-
policy = asyncio.DefaultEventLoopPolicy()
183-
if platform.system() == "Windows":
184-
policy = asyncio.WindowsProactorEventLoopPolicy() # type: ignore[attr-defined]
185-
cls._loop = policy.new_event_loop()
186-
187-
if not getattr(cls, "_thread", None):
188-
cls._thread = threading.Thread(target=cls._loop.run_forever, daemon=True)
189-
cls._thread.start()
190-
logger.info("Started loop on separate thread: %s", cls._loop)
191-
asyncio.run_coroutine_threadsafe(cls._process_queue(), cls._loop)
192-
193-
@classmethod
194-
def stop(cls, download_handler_id: int) -> None:
195-
"""Wait until all handlers are closed to stop the event loop and join the thread.
196-
Should be called from the Scrapy thread.
197-
"""
198-
cls._stop_events[download_handler_id].set()
199-
if all(ev.is_set() for ev in cls._stop_events.values()):
200-
asyncio.run_coroutine_threadsafe(cls._coro_queue.join(), cls._loop)
201-
cls._loop.call_soon_threadsafe(cls._loop.stop)
202-
cls._thread.join()
217+
log_msg = "[Context=%s] Response: <%i %s>"
218+
logger.debug(
219+
log_msg,
220+
*log_args,
221+
extra={
222+
"spider": spider,
223+
"context_name": context_name,
224+
"playwright_response_url": response.url,
225+
"playwright_response_status": response.status,
226+
},
227+
)
228+
229+
return _log_response

0 commit comments

Comments
 (0)