File
Blob: src/pyodide/internal/patches/httpx.py
| 1 | """A patch to make async httpx work using JavaScript fetch.""" |
| 2 | |
| 3 | from contextlib import contextmanager |
| 4 | |
| 5 | from httpx._client import AsyncClient, BoundAsyncStream, logger |
| 6 | from httpx._models import Headers, Request, Response |
| 7 | from httpx._transports.default import AsyncResponseStream |
| 8 | from httpx._types import AsyncByteStream |
| 9 | from httpx._utils import Timer |
| 10 | from js import Headers as js_Headers |
| 11 | from js import fetch |
| 12 | |
| 13 | from pyodide.ffi import create_proxy |
| 14 | |
| 15 | |
| 16 | @contextmanager |
| 17 | def acquire_buffer(content): |
| 18 | """Acquire a Uint8Array view of a bytes object""" |
| 19 | if not content: |
| 20 | yield None |
| 21 | return |
| 22 | body_px = create_proxy(content) |
| 23 | body_buf = body_px.getBuffer("u8") |
| 24 | try: |
| 25 | yield body_buf.data |
| 26 | finally: |
| 27 | body_px.destroy() |
| 28 | body_buf.release() |
| 29 | |
| 30 | |
| 31 | async def js_readable_stream_iter(js_readable_stream): |
| 32 | """Readable streams are supposed to be async iterators some day but they |
| 33 | aren't yet. In the meantime, this is an adaptor that produces an async |
| 34 | iterator from a readable stream. |
| 35 | """ |
| 36 | reader = js_readable_stream.getReader() |
| 37 | while True: |
| 38 | res = await reader.read() |
| 39 | if res.done: |
| 40 | return |
| 41 | b = res.value.to_bytes() |
| 42 | print("js_readable_stream_iter", b) |
| 43 | yield b |
| 44 | |
| 45 | |
| 46 | async def _send_single_request(self, request: Request) -> Response: |
| 47 | """ |
| 48 | Sends a single request, without handling any redirections. |
| 49 | |
| 50 | This is the function we're patching here... |
| 51 | """ |
| 52 | timer = Timer() |
| 53 | await timer.async_start() |
| 54 | |
| 55 | if not isinstance(request.stream, AsyncByteStream): |
| 56 | raise TypeError( |
| 57 | "Attempted to send an sync request with an AsyncClient instance." |
| 58 | ) |
| 59 | |
| 60 | # BEGIN MODIFIED PART |
| 61 | js_headers = js_Headers.new(request.headers.multi_items()) |
| 62 | with acquire_buffer(request.content) as body: |
| 63 | js_resp = await fetch( |
| 64 | str(request.url), method=request.method, headers=js_headers, body=body |
| 65 | ) |
| 66 | |
| 67 | py_headers = Headers(js_resp.headers) |
| 68 | # Unset content-encoding b/c Javascript fetch already handled unpacking. If |
| 69 | # we leave it we will get errors when httpx tries to unpack a second time. |
| 70 | py_headers.pop("content-encoding", None) |
| 71 | response = Response( |
| 72 | status_code=js_resp.status, |
| 73 | headers=py_headers, |
| 74 | stream=AsyncResponseStream(js_readable_stream_iter(js_resp.body)), |
| 75 | ) |
| 76 | # END MODIFIED PART |
| 77 | |
| 78 | assert isinstance(response.stream, AsyncByteStream) |
| 79 | response.request = request |
| 80 | response.stream = BoundAsyncStream(response.stream, response=response, timer=timer) |
| 81 | self.cookies.extract_cookies(response) |
| 82 | response.default_encoding = self._default_encoding |
| 83 | |
| 84 | logger.info( |
| 85 | 'HTTP Request: %s %s "%s %d %s"', |
| 86 | request.method, |
| 87 | request.url, |
| 88 | response.http_version, |
| 89 | response.status_code, |
| 90 | response.reason_phrase, |
| 91 | ) |
| 92 | |
| 93 | return response |
| 94 | |
| 95 | |
| 96 | AsyncClient._send_single_request = _send_single_request |