http.py 5.8 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157
  1. # -*- coding: utf-8 -*-
  2. """
  3. oss2.http
  4. ~~~~~~~~
  5. 这个模块包含了HTTP Adapters。尽管OSS Python SDK内部使用requests库进行HTTP通信,但是对使用者是透明的。
  6. 该模块中的 `Session` 、 `Request` 、`Response` 对requests的对应的类做了简单的封装。
  7. """
  8. import platform
  9. import requests
  10. from requests.structures import CaseInsensitiveDict
  11. from . import __version__, defaults
  12. from .compat import to_bytes
  13. from .exceptions import RequestError
  14. from .utils import file_object_remaining_bytes, SizedFileAdapter
  15. import logging
  16. USER_AGENT = 'aliyun-sdk-python/{0}({1}/{2}/{3};{4})'.format(
  17. __version__, platform.system(), platform.release(), platform.machine(), platform.python_version())
  18. logger = logging.getLogger(__name__)
  19. class Session(object):
  20. """属于同一个Session的请求共享一组连接池,如有可能也会重用HTTP连接。"""
  21. def __init__(self, pool_size=None, adapter=None):
  22. self.session = requests.Session()
  23. psize = pool_size or defaults.connection_pool_size
  24. if adapter is None:
  25. self.session.mount('http://', requests.adapters.HTTPAdapter(pool_connections=psize, pool_maxsize=psize))
  26. self.session.mount('https://', requests.adapters.HTTPAdapter(pool_connections=psize, pool_maxsize=psize))
  27. else:
  28. self.session.mount('http://', adapter)
  29. self.session.mount('https://', adapter)
  30. def do_request(self, req, timeout):
  31. try:
  32. logger.debug("Send request, method: {0}, url: {1}, params: {2}, headers: {3}, timeout: {4}, proxies: {5}".format(
  33. req.method, req.url, req.params, req.headers, timeout, req.proxies))
  34. return Response(self.session.request(req.method, req.url,
  35. data=req.data,
  36. params=req.params,
  37. headers=req.headers,
  38. stream=True,
  39. timeout=timeout,
  40. proxies=req.proxies))
  41. except requests.RequestException as e:
  42. raise RequestError(e)
  43. class Request(object):
  44. def __init__(self, method, url,
  45. data=None,
  46. params=None,
  47. headers=None,
  48. app_name='',
  49. proxies=None,
  50. region=None,
  51. product=None,
  52. cloudbox_id=None):
  53. self.method = method
  54. self.url = url
  55. self.data = _convert_request_body(data)
  56. self.params = params or {}
  57. self.proxies = proxies
  58. self.region = region
  59. self.product = product
  60. self.cloudbox_id = cloudbox_id
  61. if not isinstance(headers, CaseInsensitiveDict):
  62. self.headers = CaseInsensitiveDict(headers)
  63. else:
  64. self.headers = headers
  65. # tell requests not to add 'Accept-Encoding: gzip, deflate' by default
  66. if 'Accept-Encoding' not in self.headers:
  67. self.headers['Accept-Encoding'] = None
  68. if 'User-Agent' not in self.headers:
  69. if app_name:
  70. self.headers['User-Agent'] = USER_AGENT + '/' + app_name
  71. else:
  72. self.headers['User-Agent'] = USER_AGENT
  73. logger.debug("Init request, method: {0}, url: {1}, params: {2}, headers: {3}".format(method, url, params,
  74. headers))
  75. _CHUNK_SIZE = 8 * 1024
  76. class Response(object):
  77. def __init__(self, response):
  78. self.response = response
  79. self.status = response.status_code
  80. self.headers = response.headers
  81. self.request_id = response.headers.get('x-oss-request-id', '')
  82. # When a response contains no body, iter_content() cannot
  83. # be run twice (requests.exceptions.StreamConsumedError will be raised).
  84. # For details of the issue, please see issue #82
  85. #
  86. # To work around this issue, we simply return b'' when everything has been read.
  87. #
  88. # Note you cannot use self.response.raw.read() to implement self.read(), because
  89. # raw.read() does not uncompress response body when the encoding is gzip etc., and
  90. # we try to avoid depends on details of self.response.raw.
  91. self.__all_read = False
  92. logger.debug("Get response headers, req-id:{0}, status: {1}, headers: {2}".format(self.request_id, self.status,
  93. self.headers))
  94. def read(self, amt=None):
  95. if self.__all_read:
  96. return b''
  97. if amt is None:
  98. content_list = []
  99. for chunk in self.response.iter_content(_CHUNK_SIZE):
  100. content_list.append(chunk)
  101. content = b''.join(content_list)
  102. self.__all_read = True
  103. return content
  104. else:
  105. try:
  106. return next(self.response.iter_content(amt))
  107. except StopIteration:
  108. self.__all_read = True
  109. return b''
  110. def __iter__(self):
  111. return self.response.iter_content(_CHUNK_SIZE)
  112. # requests对于具有fileno()方法的file object,会用fileno()的返回值作为Content-Length。
  113. # 这对于已经读取了部分内容,或执行了seek()的file object是不正确的。
  114. #
  115. # _convert_request_body()对于支持seek()和tell() file object,确保是从
  116. # 当前位置读取,且只读取当前位置到文件结束的内容。
  117. def _convert_request_body(data):
  118. data = to_bytes(data)
  119. if hasattr(data, '__len__'):
  120. return data
  121. if hasattr(data, 'seek') and hasattr(data, 'tell'):
  122. return SizedFileAdapter(data, file_object_remaining_bytes(data))
  123. return data