codekingpro/portable-devtools
114k
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 