Team Ai
Datasetpublic

codekingpro/portable-devtools

sourceHugging Faceupdated 5mo agoView on Hugging Face
1likes14kdownloads
curl_httpclient.py591 linesDownload Raw Back to tornado
1#
2# Copyright 2009 Facebook
3#
4# Licensed under the Apache License, Version 2.0 (the "License"); you may
5# not use this file except in compliance with the License. You may obtain
6# a copy of the License at
7#
8#     http://www.apache.org/licenses/LICENSE-2.0
9#
10# Unless required by applicable law or agreed to in writing, software
11# distributed under the License is distributed on an "AS IS" BASIS, WITHOUT
12# WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the
13# License for the specific language governing permissions and limitations
14# under the License.
15
16"""Non-blocking HTTP client implementation using pycurl."""
17
18import collections
19import functools
20import logging
21import pycurl
22import re
23import threading
24import time
25from io import BytesIO
26
27from tornado import httputil
28from tornado import ioloop
29
30from tornado.escape import utf8, native_str
31from tornado.httpclient import (
32    HTTPRequest,
33    HTTPResponse,
34    HTTPError,
35    AsyncHTTPClient,
36    main,
37)
38from tornado.log import app_log
39
40from typing import Dict, Any, Callable, Union, Optional
41import typing
42
43if typing.TYPE_CHECKING:
44    from typing import Deque, Tuple  # noqa: F401
45
46curl_log = logging.getLogger("tornado.curl_httpclient")
47
48CR_OR_LF_RE = re.compile(b"\r|\n")
49
50
51class CurlAsyncHTTPClient(AsyncHTTPClient):
52    def initialize(  # type: ignore
53        self, max_clients: int = 10, defaults: Optional[Dict[str, Any]] = None
54    ) -> None:
55        super().initialize(defaults=defaults)
56        # Typeshed is incomplete for CurlMulti, so just use Any for now.
57        self._multi = pycurl.CurlMulti()  # type: Any
58        self._multi.setopt(pycurl.M_TIMERFUNCTION, self._set_timeout)
59        self._multi.setopt(pycurl.M_SOCKETFUNCTION, self._handle_socket)
60        self._curls = [self._curl_create() for i in range(max_clients)]
61        self._free_list = self._curls[:]
62        self._requests = (
63            collections.deque()
64        )  # type: Deque[Tuple[HTTPRequest, Callable[[HTTPResponse], None], float]]
65        self._fds = {}  # type: Dict[int, int]
66        self._timeout = None  # type: Optional[object]
67
68        # libcurl has bugs that sometimes cause it to not report all
69        # relevant file descriptors and timeouts to TIMERFUNCTION/
70        # SOCKETFUNCTION.  Mitigate the effects of such bugs by
71        # forcing a periodic scan of all active requests.
72        self._force_timeout_callback = ioloop.PeriodicCallback(
73            self._handle_force_timeout, 1000
74        )
75        self._force_timeout_callback.start()
76
77        # Work around a bug in libcurl 7.29.0: Some fields in the curl
78        # multi object are initialized lazily, and its destructor will
79        # segfault if it is destroyed without having been used.  Add
80        # and remove a dummy handle to make sure everything is
81        # initialized.
82        dummy_curl_handle = pycurl.Curl()
83        self._multi.add_handle(dummy_curl_handle)
84        self._multi.remove_handle(dummy_curl_handle)
85
86    def close(self) -> None:
87        self._force_timeout_callback.stop()
88        if self._timeout is not None:
89            self.io_loop.remove_timeout(self._timeout)
90        for curl in self._curls:
91            curl.close()
92        self._multi.close()
93        super().close()
94
95        # Set below properties to None to reduce the reference count of current
96        # instance, because those properties hold some methods of current
97        # instance that will case circular reference.
98        self._force_timeout_callback = None  # type: ignore
99        self._multi = None
100
101    def fetch_impl(
102        self, request: HTTPRequest, callback: Callable[[HTTPResponse], None]
103    ) -> None:
104        self._requests.append((request, callback, self.io_loop.time()))
105        self._process_queue()
106        self._set_timeout(0)
107
108    def _handle_socket(self, event: int, fd: int, multi: Any, data: bytes) -> None:
109        """Called by libcurl when it wants to change the file descriptors
110        it cares about.
111        """
112        event_map = {
113            pycurl.POLL_NONE: ioloop.IOLoop.NONE,
114            pycurl.POLL_IN: ioloop.IOLoop.READ,
115            pycurl.POLL_OUT: ioloop.IOLoop.WRITE,
116            pycurl.POLL_INOUT: ioloop.IOLoop.READ | ioloop.IOLoop.WRITE,
117        }
118        if event == pycurl.POLL_REMOVE:
119            if fd in self._fds:
120                self.io_loop.remove_handler(fd)
121                del self._fds[fd]
122        else:
123            ioloop_event = event_map[event]
124            # libcurl sometimes closes a socket and then opens a new
125            # one using the same FD without giving us a POLL_NONE in
126            # between.  This is a problem with the epoll IOLoop,
127            # because the kernel can tell when a socket is closed and
128            # removes it from the epoll automatically, causing future
129            # update_handler calls to fail.  Since we can't tell when
130            # this has happened, always use remove and re-add
131            # instead of update.
132            if fd in self._fds:
133                self.io_loop.remove_handler(fd)
134            self.io_loop.add_handler(fd, self._handle_events, ioloop_event)
135            self._fds[fd] = ioloop_event
136
137    def _set_timeout(self, msecs: int) -> None:
138        """Called by libcurl to schedule a timeout."""
139        if self._timeout is not None:
140            self.io_loop.remove_timeout(self._timeout)
141        self._timeout = self.io_loop.add_timeout(
142            self.io_loop.time() + msecs / 1000.0, self._handle_timeout
143        )
144
145    def _handle_events(self, fd: int, events: int) -> None:
146        """Called by IOLoop when there is activity on one of our
147        file descriptors.
148        """
149        action = 0
150        if events & ioloop.IOLoop.READ:
151            action |= pycurl.CSELECT_IN
152        if events & ioloop.IOLoop.WRITE:
153            action |= pycurl.CSELECT_OUT
154        while True:
155            try:
156                ret, num_handles = self._multi.socket_action(fd, action)
157            except pycurl.error as e:
158                ret = e.args[0]
159            if ret != pycurl.E_CALL_MULTI_PERFORM:
160                break
161        self._finish_pending_requests()
162
163    def _handle_timeout(self) -> None:
164        """Called by IOLoop when the requested timeout has passed."""
165        self._timeout = None
166        while True:
167            try:
168                ret, num_handles = self._multi.socket_action(pycurl.SOCKET_TIMEOUT, 0)
169            except pycurl.error as e:
170                ret = e.args[0]
171            if ret != pycurl.E_CALL_MULTI_PERFORM:
172                break
173        self._finish_pending_requests()
174
175        # In theory, we shouldn't have to do this because curl will
176        # call _set_timeout whenever the timeout changes.  However,
177        # sometimes after _handle_timeout we will need to reschedule
178        # immediately even though nothing has changed from curl's
179        # perspective.  This is because when socket_action is
180        # called with SOCKET_TIMEOUT, libcurl decides internally which
181        # timeouts need to be processed by using a monotonic clock
182        # (where available) while tornado uses python's time.time()
183        # to decide when timeouts have occurred.  When those clocks
184        # disagree on elapsed time (as they will whenever there is an
185        # NTP adjustment), tornado might call _handle_timeout before
186        # libcurl is ready.  After each timeout, resync the scheduled
187        # timeout with libcurl's current state.
188        new_timeout = self._multi.timeout()
189        if new_timeout >= 0:
190            self._set_timeout(new_timeout)
191
192    def _handle_force_timeout(self) -> None:
193        """Called by IOLoop periodically to ask libcurl to process any
194        events it may have forgotten about.
195        """
196        while True:
197            try:
198                ret, num_handles = self._multi.socket_all()
199            except pycurl.error as e:
200                ret = e.args[0]
201            if ret != pycurl.E_CALL_MULTI_PERFORM:
202                break
203        self._finish_pending_requests()
204
205    def _finish_pending_requests(self) -> None:
206        """Process any requests that were completed by the last
207        call to multi.socket_action.
208        """
209        while True:
210            num_q, ok_list, err_list = self._multi.info_read()
211            for curl in ok_list:
212                self._finish(curl)
213            for curl, errnum, errmsg in err_list:
214                self._finish(curl, errnum, errmsg)
215            if num_q == 0:
216                break
217        self._process_queue()
218
219    def _process_queue(self) -> None:
220        while True:
221            started = 0
222            while self._free_list and self._requests:
223                started += 1
224                curl = self._free_list.pop()
225                (request, callback, queue_start_time) = self._requests.popleft()
226                # TODO: Don't smuggle extra data on an attribute of the Curl object.
227                curl.info = {  # type: ignore
228                    "headers": httputil.HTTPHeaders(),
229                    "buffer": BytesIO(),
230                    "request": request,
231                    "callback": callback,
232                    "queue_start_time": queue_start_time,
233                    "curl_start_time": time.time(),
234                    "curl_start_ioloop_time": self.io_loop.current().time(),  # type: ignore
235                }
236                try:
237                    self._curl_setup_request(
238                        curl,
239                        request,
240                        curl.info["buffer"],  # type: ignore
241                        curl.info["headers"],  # type: ignore
242                    )
243                except Exception as e:
244                    # If there was an error in setup, pass it on
245                    # to the callback. Note that allowing the
246                    # error to escape here will appear to work
247                    # most of the time since we are still in the
248                    # caller's original stack frame, but when
249                    # _process_queue() is called from
250                    # _finish_pending_requests the exceptions have
251                    # nowhere to go.
252                    self._free_list.append(curl)
253                    callback(HTTPResponse(request=request, code=599, error=e))
254                else:
255                    self._multi.add_handle(curl)
256
257            if not started:
258                break
259
260    def _finish(
261        self,
262        curl: pycurl.Curl,
263        curl_error: Optional[int] = None,
264        curl_message: Optional[str] = None,
265    ) -> None:
266        info = curl.info  # type: ignore
267        curl.info = None  # type: ignore
268        self._multi.remove_handle(curl)
269        self._free_list.append(curl)
270        buffer = info["buffer"]
271        if curl_error:
272            assert curl_message is not None
273            error = CurlError(curl_error, curl_message)  # type: Optional[CurlError]
274            assert error is not None
275            code = error.code
276            effective_url = None
277            buffer.close()
278            buffer = None
279        else:
280            error = None
281            code = curl.getinfo(pycurl.HTTP_CODE)
282            effective_url = curl.getinfo(pycurl.EFFECTIVE_URL)
283            buffer.seek(0)
284        # the various curl timings are documented at
285        # http://curl.haxx.se/libcurl/c/curl_easy_getinfo.html
286        time_info = dict(
287            queue=info["curl_start_ioloop_time"] - info["queue_start_time"],
288            namelookup=curl.getinfo(pycurl.NAMELOOKUP_TIME),
289            connect=curl.getinfo(pycurl.CONNECT_TIME),
290            appconnect=curl.getinfo(pycurl.APPCONNECT_TIME),
291            pretransfer=curl.getinfo(pycurl.PRETRANSFER_TIME),
292            starttransfer=curl.getinfo(pycurl.STARTTRANSFER_TIME),
293            total=curl.getinfo(pycurl.TOTAL_TIME),
294            redirect=curl.getinfo(pycurl.REDIRECT_TIME),
295        )
296        try:
297            info["callback"](
298                HTTPResponse(
299                    request=info["request"],
300                    code=code,
301                    headers=info["headers"],
302                    buffer=buffer,
303                    effective_url=effective_url,
304                    error=error,
305                    reason=info["headers"].get("X-Http-Reason", None),
306                    request_time=self.io_loop.time() - info["curl_start_ioloop_time"],
307                    start_time=info["curl_start_time"],
308                    time_info=time_info,
309                )
310            )
311        except Exception:
312            self.handle_callback_exception(info["callback"])
313
314    def handle_callback_exception(self, callback: Any) -> None:
315        app_log.error("Exception in callback %r", callback, exc_info=True)
316
317    def _curl_create(self) -> pycurl.Curl:
318        curl = pycurl.Curl()
319        if curl_log.isEnabledFor(logging.DEBUG):
320            curl.setopt(pycurl.VERBOSE, 1)
321            curl.setopt(pycurl.DEBUGFUNCTION, self._curl_debug)
322        if hasattr(
323            pycurl, "PROTOCOLS"
324        ):  # PROTOCOLS first appeared in pycurl 7.19.5 (2014-07-12)
325            curl.setopt(pycurl.PROTOCOLS, pycurl.PROTO_HTTP | pycurl.PROTO_HTTPS)
326            curl.setopt(pycurl.REDIR_PROTOCOLS, pycurl.PROTO_HTTP | pycurl.PROTO_HTTPS)
327        return curl
328
329    def _curl_setup_request(
330        self,
331        curl: pycurl.Curl,
332        request: HTTPRequest,
333        buffer: BytesIO,
334        headers: httputil.HTTPHeaders,
335    ) -> None:
336        curl.setopt(pycurl.URL, native_str(request.url))
337
338        # libcurl's magic "Expect: 100-continue" behavior causes delays
339        # with servers that don't support it (which include, among others,
340        # Google's OpenID endpoint).  Additionally, this behavior has
341        # a bug in conjunction with the curl_multi_socket_action API
342        # (https://sourceforge.net/tracker/?func=detail&atid=100976&aid=3039744&group_id=976),
343        # which increases the delays.  It's more trouble than it's worth,
344        # so just turn off the feature (yes, setting Expect: to an empty
345        # value is the official way to disable this)
346        if "Expect" not in request.headers:
347            request.headers["Expect"] = ""
348
349        # libcurl adds Pragma: no-cache by default; disable that too
350        if "Pragma" not in request.headers:
351            request.headers["Pragma"] = ""
352
353        encoded_headers = [
354            b"%s: %s"
355            % (native_str(k).encode("ASCII"), native_str(v).encode("ISO8859-1"))
356            for k, v in request.headers.get_all()
357        ]
358        for line in encoded_headers:
359            if CR_OR_LF_RE.search(line):
360                raise ValueError("Illegal characters in header (CR or LF): %r" % line)
361        curl.setopt(pycurl.HTTPHEADER, encoded_headers)
362
363        curl.setopt(
364            pycurl.HEADERFUNCTION,
365            functools.partial(
366                self._curl_header_callback, headers, request.header_callback
367            ),
368        )
369        if request.streaming_callback:
370
371            def write_function(b: Union[bytes, bytearray]) -> int:
372                assert request.streaming_callback is not None
373                self.io_loop.add_callback(request.streaming_callback, b)
374                return len(b)
375
376        else:
377            write_function = buffer.write  # type: ignore
378        curl.setopt(pycurl.WRITEFUNCTION, write_function)
379        curl.setopt(pycurl.FOLLOWLOCATION, request.follow_redirects)
380        curl.setopt(pycurl.MAXREDIRS, request.max_redirects)
381        assert request.connect_timeout is not None
382        curl.setopt(pycurl.CONNECTTIMEOUT_MS, int(1000 * request.connect_timeout))
383        assert request.request_timeout is not None
384        curl.setopt(pycurl.TIMEOUT_MS, int(1000 * request.request_timeout))
385        if request.user_agent:
386            curl.setopt(pycurl.USERAGENT, native_str(request.user_agent))
387        else:
388            curl.setopt(pycurl.USERAGENT, "Mozilla/5.0 (compatible; pycurl)")
389        if request.network_interface:
390            curl.setopt(pycurl.INTERFACE, request.network_interface)
391        if request.decompress_response:
392            curl.setopt(pycurl.ENCODING, "gzip,deflate")
393        else:
394            curl.setopt(pycurl.ENCODING, None)
395        if request.proxy_host and request.proxy_port:
396            curl.setopt(pycurl.PROXY, request.proxy_host)
397            curl.setopt(pycurl.PROXYPORT, request.proxy_port)
398            if request.proxy_username:
399                assert request.proxy_password is not None
400                credentials = httputil.encode_username_password(
401                    request.proxy_username, request.proxy_password
402                )
403                curl.setopt(pycurl.PROXYUSERPWD, credentials)
404
405            if request.proxy_auth_mode is None or request.proxy_auth_mode == "basic":
406                curl.setopt(pycurl.PROXYAUTH, pycurl.HTTPAUTH_BASIC)
407            elif request.proxy_auth_mode == "digest":
408                curl.setopt(pycurl.PROXYAUTH, pycurl.HTTPAUTH_DIGEST)
409            else:
410                raise ValueError(
411                    "Unsupported proxy_auth_mode %s" % request.proxy_auth_mode
412                )
413        else:
414            try:
415                curl.unsetopt(pycurl.PROXY)
416            except TypeError:  # not supported, disable proxy
417                curl.setopt(pycurl.PROXY, "")
418            curl.unsetopt(pycurl.PROXYUSERPWD)
419        if request.validate_cert:
420            curl.setopt(pycurl.SSL_VERIFYPEER, 1)
421            curl.setopt(pycurl.SSL_VERIFYHOST, 2)
422        else:
423            curl.setopt(pycurl.SSL_VERIFYPEER, 0)
424            curl.setopt(pycurl.SSL_VERIFYHOST, 0)
425        if request.ca_certs is not None:
426            curl.setopt(pycurl.CAINFO, request.ca_certs)
427        else:
428            # There is no way to restore pycurl.CAINFO to its default value
429            # (Using unsetopt makes it reject all certificates).
430            # I don't see any way to read the default value from python so it
431            # can be restored later.  We'll have to just leave CAINFO untouched
432            # if no ca_certs file was specified, and require that if any
433            # request uses a custom ca_certs file, they all must.
434            pass
435
436        if request.allow_ipv6 is False:
437            # Curl behaves reasonably when DNS resolution gives an ipv6 address
438            # that we can't reach, so allow ipv6 unless the user asks to disable.
439            curl.setopt(pycurl.IPRESOLVE, pycurl.IPRESOLVE_V4)
440        else:
441            curl.setopt(pycurl.IPRESOLVE, pycurl.IPRESOLVE_WHATEVER)
442
443        # Set the request method through curl's irritating interface which makes
444        # up names for almost every single method
445        curl_options = {
446            "GET": pycurl.HTTPGET,
447            "POST": pycurl.POST,
448            "PUT": pycurl.UPLOAD,
449            "HEAD": pycurl.NOBODY,
450        }
451        custom_methods = {"DELETE", "OPTIONS", "PATCH"}
452        for o in curl_options.values():
453            curl.setopt(o, False)
454        if request.method in curl_options:
455            curl.unsetopt(pycurl.CUSTOMREQUEST)
456            curl.setopt(curl_options[request.method], True)
457        elif request.allow_nonstandard_methods or request.method in custom_methods:
458            curl.setopt(pycurl.CUSTOMREQUEST, request.method)
459        else:
460            raise KeyError("unknown method " + request.method)
461
462        body_expected = request.method in ("POST", "PATCH", "PUT")
463        body_present = request.body is not None
464        if not request.allow_nonstandard_methods:
465            # Some HTTP methods nearly always have bodies while others
466            # almost never do. Fail in this case unless the user has
467            # opted out of sanity checks with allow_nonstandard_methods.
468            if (body_expected and not body_present) or (
469                body_present and not body_expected
470            ):
471                raise ValueError(
472                    "Body must %sbe None for method %s (unless "
473                    "allow_nonstandard_methods is true)"
474                    % ("not " if body_expected else "", request.method)
475                )
476
477        if body_expected or body_present:
478            if request.method == "GET":
479                # Even with `allow_nonstandard_methods` we disallow
480                # GET with a body (because libcurl doesn't allow it
481                # unless we use CUSTOMREQUEST). While the spec doesn't
482                # forbid clients from sending a body, it arguably
483                # disallows the server from doing anything with them.
484                raise ValueError("Body must be None for GET request")
485            request_buffer = BytesIO(utf8(request.body or ""))
486
487            def ioctl(cmd: int) -> None:
488                if cmd == curl.IOCMD_RESTARTREAD:  # type: ignore
489                    request_buffer.seek(0)
490
491            curl.setopt(pycurl.READFUNCTION, request_buffer.read)
492            curl.setopt(pycurl.IOCTLFUNCTION, ioctl)
493            if request.method == "POST":
494                curl.setopt(pycurl.POSTFIELDSIZE, len(request.body or ""))
495            else:
496                curl.setopt(pycurl.UPLOAD, True)
497                curl.setopt(pycurl.INFILESIZE, len(request.body or ""))
498
499        if request.auth_username is not None:
500            assert request.auth_password is not None
501            if request.auth_mode is None or request.auth_mode == "basic":
502                curl.setopt(pycurl.HTTPAUTH, pycurl.HTTPAUTH_BASIC)
503            elif request.auth_mode == "digest":
504                curl.setopt(pycurl.HTTPAUTH, pycurl.HTTPAUTH_DIGEST)
505            else:
506                raise ValueError("Unsupported auth_mode %s" % request.auth_mode)
507
508            userpwd = httputil.encode_username_password(
509                request.auth_username, request.auth_password
510            )
511            curl.setopt(pycurl.USERPWD, userpwd)
512            curl_log.debug(
513                "%s %s (username: %r)",
514                request.method,
515                request.url,
516                request.auth_username,
517            )
518        else:
519            curl.unsetopt(pycurl.USERPWD)
520            curl_log.debug("%s %s", request.method, request.url)
521
522        if request.client_cert is not None:
523            curl.setopt(pycurl.SSLCERT, request.client_cert)
524
525        if request.client_key is not None:
526            curl.setopt(pycurl.SSLKEY, request.client_key)
527
528        if request.ssl_options is not None:
529            raise ValueError("ssl_options not supported in curl_httpclient")
530
531        if threading.active_count() > 1:
532            # libcurl/pycurl is not thread-safe by default.  When multiple threads
533            # are used, signals should be disabled.  This has the side effect
534            # of disabling DNS timeouts in some environments (when libcurl is
535            # not linked against ares), so we don't do it when there is only one
536            # thread.  Applications that use many short-lived threads may need
537            # to set NOSIGNAL manually in a prepare_curl_callback since
538            # there may not be any other threads running at the time we call
539            # threading.activeCount.
540            curl.setopt(pycurl.NOSIGNAL, 1)
541        if request.prepare_curl_callback is not None:
542            request.prepare_curl_callback(curl)
543
544    def _curl_header_callback(
545        self,
546        headers: httputil.HTTPHeaders,
547        header_callback: Callable[[str], None],
548        header_line_bytes: bytes,
549    ) -> None:
550        header_line = native_str(header_line_bytes.decode("latin1"))
551        if header_callback is not None:
552            self.io_loop.add_callback(header_callback, header_line)
553        # header_line as returned by curl includes the end-of-line characters.
554        # whitespace at the start should be preserved to allow multi-line headers
555        header_line = header_line.rstrip()
556        if header_line.startswith("HTTP/"):
557            headers.clear()
558            try:
559                (_version, _code, reason) = httputil.parse_response_start_line(
560                    header_line
561                )
562                header_line = "X-Http-Reason: %s" % reason
563            except httputil.HTTPInputError:
564                return
565        if not header_line:
566            return
567        headers.parse_line(header_line)
568
569    def _curl_debug(self, debug_type: int, debug_msg: str) -> None:
570        debug_types = ("I", "<", ">", "<", ">")
571        if debug_type == 0:
572            debug_msg = native_str(debug_msg)
573            curl_log.debug("%s", debug_msg.strip())
574        elif debug_type in (1, 2):
575            debug_msg = native_str(debug_msg)
576            for line in debug_msg.splitlines():
577                curl_log.debug("%s %s", debug_types[debug_type], line)
578        elif debug_type == 4:
579            curl_log.debug("%s %r", debug_types[debug_type], debug_msg)
580
581
582class CurlError(HTTPError):
583    def __init__(self, errno: int, message: str) -> None:
584        HTTPError.__init__(self, 599, message)
585        self.errno = errno
586
587
588if __name__ == "__main__":
589    AsyncHTTPClient.configure(CurlAsyncHTTPClient)
590    main()
591 
codekingpro/portable-devtools · Team Ai