From cebc873cd0f443c5b04e5852bf4f982126617e82 Mon Sep 17 00:00:00 2001 From: Xiang Yan Date: Thu, 11 Mar 2021 09:33:36 -0800 Subject: [PATCH 01/36] add raw streaming support --- sdk/core/azure-core/CLIENT_LIBRARY_DEVELOPER.md | 2 +- .../azure/core/pipeline/transport/_aiohttp.py | 10 +++++----- .../azure/core/pipeline/transport/_base.py | 4 ++-- .../azure/core/pipeline/transport/_base_async.py | 2 +- .../core/pipeline/transport/_requests_asyncio.py | 10 +++++----- .../core/pipeline/transport/_requests_basic.py | 15 ++++++++++----- .../core/pipeline/transport/_requests_trio.py | 8 +++++--- .../azure-core/tests/test_stream_generator.py | 10 ++++++++++ 8 files changed, 39 insertions(+), 22 deletions(-) diff --git a/sdk/core/azure-core/CLIENT_LIBRARY_DEVELOPER.md b/sdk/core/azure-core/CLIENT_LIBRARY_DEVELOPER.md index f9556d072e46..7a29caae3e08 100644 --- a/sdk/core/azure-core/CLIENT_LIBRARY_DEVELOPER.md +++ b/sdk/core/azure-core/CLIENT_LIBRARY_DEVELOPER.md @@ -269,7 +269,7 @@ class HttpResponse(object): def text(self, encoding=None): """Return the whole body as a string.""" - def stream_download(self, chunk_size=None, callback=None): + def stream_download(self, pipeline, raw=False): """Generator for streaming request body data. Should be implemented by sub-classes if streaming download is supported. diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py b/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py index af553bc71b30..bfc6f6d470c8 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py @@ -197,16 +197,16 @@ class AioHttpStreamDownloadGenerator(AsyncIterator): :param pipeline: The pipeline object :param response: The client response object. - :param block_size: block size of data sent over connection. - :type block_size: int + :param raw: If returns the raw stream. """ - def __init__(self, pipeline: Pipeline, response: AsyncHttpResponse) -> None: + def __init__(self, pipeline: Pipeline, response: AsyncHttpResponse, raw: bool=False) -> None: self.pipeline = pipeline self.request = response.request self.response = response self.block_size = response.block_size self.content_length = int(response.internal_response.headers.get('Content-Length', 0)) self.downloaded = 0 + self._raw = raw def __len__(self): return self.content_length @@ -269,13 +269,13 @@ async def load_body(self) -> None: """Load in memory the body, so it could be accessible from sync methods.""" self._body = await self.internal_response.read() - def stream_download(self, pipeline) -> AsyncIteratorType[bytes]: + def stream_download(self, pipeline, raw=True) -> AsyncIteratorType[bytes]: """Generator for streaming response body data. :param pipeline: The pipeline object :type pipeline: azure.core.pipeline """ - return AioHttpStreamDownloadGenerator(pipeline, self) + return AioHttpStreamDownloadGenerator(pipeline, self, raw=raw) def __getstate__(self): # Be sure body is loaded in memory, otherwise not pickable and let it throw diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_base.py b/sdk/core/azure-core/azure/core/pipeline/transport/_base.py index 6531ea5179f1..8b4304bdea4c 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_base.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_base.py @@ -605,8 +605,8 @@ def __repr__(self): class HttpResponse(_HttpResponseBase): # pylint: disable=abstract-method - def stream_download(self, pipeline): - # type: (PipelineType) -> Iterator[bytes] + def stream_download(self, pipeline, raw=False): + # type: (PipelineType, bool) -> Iterator[bytes] """Generator for streaming request body data. Should be implemented by sub-classes if streaming download diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_base_async.py b/sdk/core/azure-core/azure/core/pipeline/transport/_base_async.py index bfc51ef6109b..6dab35f1e992 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_base_async.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_base_async.py @@ -124,7 +124,7 @@ class AsyncHttpResponse(_HttpResponseBase): # pylint: disable=abstract-method Allows for the asynchronous streaming of data from the response. """ - def stream_download(self, pipeline) -> AsyncIteratorType[bytes]: + def stream_download(self, pipeline, raw=False) -> AsyncIteratorType[bytes]: """Generator for streaming response body data. Should be implemented by sub-classes if streaming download diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_asyncio.py b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_asyncio.py index 90c53675f866..b414d0db46cf 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_asyncio.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_asyncio.py @@ -138,10 +138,9 @@ class AsyncioStreamDownloadGenerator(AsyncIterator): :param pipeline: The pipeline object :param response: The response object. - :param generator iter_content_func: Iterator for response data. - :param int content_length: size of body in bytes. + :param raw: If returns the raw stream. """ - def __init__(self, pipeline: Pipeline, response: AsyncHttpResponse) -> None: + def __init__(self, pipeline: Pipeline, response: AsyncHttpResponse, raw: bool=False) -> None: self.pipeline = pipeline self.request = response.request self.response = response @@ -149,6 +148,7 @@ def __init__(self, pipeline: Pipeline, response: AsyncHttpResponse) -> None: self.iter_content_func = self.response.internal_response.iter_content(self.block_size) self.content_length = int(response.headers.get('Content-Length', 0)) self.downloaded = 0 + self._raw = raw def __len__(self): return self.content_length @@ -178,6 +178,6 @@ async def __anext__(self): class AsyncioRequestsTransportResponse(AsyncHttpResponse, RequestsTransportResponse): # type: ignore """Asynchronous streaming of data from the response. """ - def stream_download(self, pipeline) -> AsyncIteratorType[bytes]: # type: ignore + def stream_download(self, pipeline, raw=False) -> AsyncIteratorType[bytes]: # type: ignore """Generator for streaming request body data.""" - return AsyncioStreamDownloadGenerator(pipeline, self) # type: ignore + return AsyncioStreamDownloadGenerator(pipeline, self, raw=raw) # type: ignore diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_basic.py b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_basic.py index ed8a7382c55d..54c716597de5 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_basic.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_basic.py @@ -98,14 +98,16 @@ class StreamDownloadGenerator(object): :param pipeline: The pipeline object :param response: The response object. + :param raw: If returns the raw stream. """ - def __init__(self, pipeline, response): + def __init__(self, pipeline, response, raw=False): self.pipeline = pipeline self.request = response.request self.response = response self.block_size = response.block_size self.iter_content_func = self.response.internal_response.iter_content(self.block_size) self.content_length = int(response.headers.get('Content-Length', 0)) + self._raw = raw def __len__(self): return self.content_length @@ -115,7 +117,10 @@ def __iter__(self): def __next__(self): try: - chunk = next(self.iter_content_func) + if self._raw: + chunk = self.response.internal_response.raw.read(self.block_size, decode_content=False) + else: + chunk = next(self.iter_content_func) if not chunk: raise StopIteration() return chunk @@ -134,10 +139,10 @@ def __next__(self): class RequestsTransportResponse(HttpResponse, _RequestsTransportResponseBase): """Streaming of data from the response. """ - def stream_download(self, pipeline): - # type: (PipelineType) -> Iterator[bytes] + def stream_download(self, pipeline, raw=True): + # type: (PipelineType, bool) -> Iterator[bytes] """Generator for streaming request body data.""" - return StreamDownloadGenerator(pipeline, self) + return StreamDownloadGenerator(pipeline, self, raw=raw) class RequestsTransport(HttpTransport): diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_trio.py b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_trio.py index 04ddd453bbf5..3e140ccfed9c 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_trio.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_trio.py @@ -54,8 +54,9 @@ class TrioStreamDownloadGenerator(AsyncIterator): :param pipeline: The pipeline object :param response: The response object. + :param raw: If returns the raw stream. """ - def __init__(self, pipeline: Pipeline, response: AsyncHttpResponse) -> None: + def __init__(self, pipeline: Pipeline, response: AsyncHttpResponse, raw: bool=False) -> None: self.pipeline = pipeline self.request = response.request self.response = response @@ -63,6 +64,7 @@ def __init__(self, pipeline: Pipeline, response: AsyncHttpResponse) -> None: self.iter_content_func = self.response.internal_response.iter_content(self.block_size) self.content_length = int(response.headers.get('Content-Length', 0)) self.downloaded = 0 + self._raw = raw def __len__(self): return self.content_length @@ -95,10 +97,10 @@ async def __anext__(self): class TrioRequestsTransportResponse(AsyncHttpResponse, RequestsTransportResponse): # type: ignore """Asynchronous streaming of data from the response. """ - def stream_download(self, pipeline) -> AsyncIteratorType[bytes]: # type: ignore + def stream_download(self, pipeline, raw=True) -> AsyncIteratorType[bytes]: # type: ignore """Generator for streaming response data. """ - return TrioStreamDownloadGenerator(pipeline, self) + return TrioStreamDownloadGenerator(pipeline, self, raw=raw) class TrioRequestsTransport(RequestsAsyncTransportBase): # type: ignore diff --git a/sdk/core/azure-core/tests/test_stream_generator.py b/sdk/core/azure-core/tests/test_stream_generator.py index 8b30b00a4f50..d4ad647b29c6 100644 --- a/sdk/core/azure-core/tests/test_stream_generator.py +++ b/sdk/core/azure-core/tests/test_stream_generator.py @@ -80,6 +80,16 @@ def stream(self, chunk_size, decode_content=False): left -= len(data) yield data + def read(self, chunk_size, decode_content=False): + assert chunk_size == block_size + left = total_response_size + while left > 0: + if left <= block_size: + raise requests.exceptions.ConnectionError() + data = b"X" * min(chunk_size, left) + left -= len(data) + yield data + def close(self): pass From 15d2f4dae08287c21b441d064847cfc324d3eede Mon Sep 17 00:00:00 2001 From: Xiang Yan Date: Wed, 14 Apr 2021 13:08:47 -0700 Subject: [PATCH 02/36] update --- .../azure/core/pipeline/transport/_aiohttp.py | 5 +++-- .../pipeline/transport/_requests_asyncio.py | 22 +++++++++++++------ .../pipeline/transport/_requests_basic.py | 9 ++++---- .../core/pipeline/transport/_requests_trio.py | 5 +++-- .../test_stream_generator_async.py | 11 ++++++++++ .../azure-core/tests/test_stream_generator.py | 14 +++++++----- 6 files changed, 44 insertions(+), 22 deletions(-) diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py b/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py index bfc6f6d470c8..6c951d5ba07b 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py @@ -204,9 +204,10 @@ def __init__(self, pipeline: Pipeline, response: AsyncHttpResponse, raw: bool=Fa self.request = response.request self.response = response self.block_size = response.block_size - self.content_length = int(response.internal_response.headers.get('Content-Length', 0)) - self.downloaded = 0 self._raw = raw + self.content_length = int(response.internal_response.headers.get('Content-Length', 0)) + if self._raw and hasattr(self.response.internal_response.raw, 'stream'): + delattr(self.response.internal_response.raw.__class__, 'stream') def __len__(self): return self.content_length diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_asyncio.py b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_asyncio.py index b414d0db46cf..46678e4ff040 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_asyncio.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_asyncio.py @@ -145,10 +145,11 @@ def __init__(self, pipeline: Pipeline, response: AsyncHttpResponse, raw: bool=Fa self.request = response.request self.response = response self.block_size = response.block_size + self._raw = raw self.iter_content_func = self.response.internal_response.iter_content(self.block_size) + if self._raw and hasattr(self.response.internal_response.raw, 'stream'): + delattr(self.response.internal_response.raw.__class__, 'stream') self.content_length = int(response.headers.get('Content-Length', 0)) - self.downloaded = 0 - self._raw = raw def __len__(self): return self.content_length @@ -156,10 +157,17 @@ def __len__(self): async def __anext__(self): loop = _get_running_loop() try: - chunk = await loop.run_in_executor( - None, - _iterate_response_content, - self.iter_content_func, + if self._raw: + chunk = await loop.run_in_executor( + None, + _iterate_response_content, + self.iter_content_func, + ) + else: + chunk = await loop.run_in_executor( + None, + _iterate_response_content, + self.iter_content_func, ) if not chunk: raise _ResponseStopIteration() @@ -178,6 +186,6 @@ async def __anext__(self): class AsyncioRequestsTransportResponse(AsyncHttpResponse, RequestsTransportResponse): # type: ignore """Asynchronous streaming of data from the response. """ - def stream_download(self, pipeline, raw=False) -> AsyncIteratorType[bytes]: # type: ignore + def stream_download(self, pipeline, raw=True) -> AsyncIteratorType[bytes]: # type: ignore """Generator for streaming request body data.""" return AsyncioStreamDownloadGenerator(pipeline, self, raw=raw) # type: ignore diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_basic.py b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_basic.py index 54c716597de5..92886d6f7c53 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_basic.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_basic.py @@ -105,9 +105,11 @@ def __init__(self, pipeline, response, raw=False): self.request = response.request self.response = response self.block_size = response.block_size + self._raw = raw self.iter_content_func = self.response.internal_response.iter_content(self.block_size) + if self._raw and hasattr(self.response.internal_response.raw, 'stream'): + delattr(self.response.internal_response.raw.__class__, 'stream') self.content_length = int(response.headers.get('Content-Length', 0)) - self._raw = raw def __len__(self): return self.content_length @@ -117,10 +119,7 @@ def __iter__(self): def __next__(self): try: - if self._raw: - chunk = self.response.internal_response.raw.read(self.block_size, decode_content=False) - else: - chunk = next(self.iter_content_func) + chunk = next(self.iter_content_func) if not chunk: raise StopIteration() return chunk diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_trio.py b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_trio.py index 3e140ccfed9c..eb0f2f79a65d 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_trio.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_trio.py @@ -61,10 +61,11 @@ def __init__(self, pipeline: Pipeline, response: AsyncHttpResponse, raw: bool=Fa self.request = response.request self.response = response self.block_size = response.block_size + self._raw = raw self.iter_content_func = self.response.internal_response.iter_content(self.block_size) + if self._raw and hasattr(self.response.internal_response.raw, 'stream'): + delattr(self.response.internal_response.raw.__class__, 'stream') self.content_length = int(response.headers.get('Content-Length', 0)) - self.downloaded = 0 - self._raw = raw def __len__(self): return self.content_length diff --git a/sdk/core/azure-core/tests/async_tests/test_stream_generator_async.py b/sdk/core/azure-core/tests/async_tests/test_stream_generator_async.py index 0a4d6017e6c4..af53b36a0045 100644 --- a/sdk/core/azure-core/tests/async_tests/test_stream_generator_async.py +++ b/sdk/core/azure-core/tests/async_tests/test_stream_generator_async.py @@ -75,6 +75,8 @@ async def test_response_streaming_error_behavior(): class FakeStreamWithConnectionError: # fake object for urllib3.response.HTTPResponse + def __init__(self): + self.total_response_size = 500 def stream(self, chunk_size, decode_content=False): assert chunk_size == block_size @@ -86,6 +88,15 @@ def stream(self, chunk_size, decode_content=False): left -= len(data) yield data + def read(self, chunk_size, decode_content=False): + assert chunk_size == block_size + if self.total_response_size > 0: + if self.total_response_size <= block_size: + raise requests.exceptions.ConnectionError() + data = b"X" * min(chunk_size, self.total_response_size) + self.total_response_size -= len(data) + return data + def close(self): pass diff --git a/sdk/core/azure-core/tests/test_stream_generator.py b/sdk/core/azure-core/tests/test_stream_generator.py index d4ad647b29c6..d3834c1b5e34 100644 --- a/sdk/core/azure-core/tests/test_stream_generator.py +++ b/sdk/core/azure-core/tests/test_stream_generator.py @@ -69,6 +69,8 @@ def test_response_streaming_error_behavior(): class FakeStreamWithConnectionError: # fake object for urllib3.response.HTTPResponse + def __init__(self): + self.total_response_size = 500 def stream(self, chunk_size, decode_content=False): assert chunk_size == block_size @@ -82,17 +84,17 @@ def stream(self, chunk_size, decode_content=False): def read(self, chunk_size, decode_content=False): assert chunk_size == block_size - left = total_response_size - while left > 0: - if left <= block_size: + if self.total_response_size > 0: + if self.total_response_size <= block_size: raise requests.exceptions.ConnectionError() - data = b"X" * min(chunk_size, left) - left -= len(data) - yield data + data = b"X" * min(chunk_size, self.total_response_size) + self.total_response_size -= len(data) + return data def close(self): pass + s = FakeStreamWithConnectionError() req_response.raw = FakeStreamWithConnectionError() response = RequestsTransportResponse( From 714a64acfee80ad4f93619e62aca0b95f705efdf Mon Sep 17 00:00:00 2001 From: Xiang Yan Date: Wed, 14 Apr 2021 13:22:08 -0700 Subject: [PATCH 03/36] update --- sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py | 3 ++- .../azure/core/pipeline/transport/_requests_asyncio.py | 3 ++- .../azure/core/pipeline/transport/_requests_basic.py | 3 ++- .../azure-core/azure/core/pipeline/transport/_requests_trio.py | 3 ++- 4 files changed, 8 insertions(+), 4 deletions(-) diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py b/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py index 6c951d5ba07b..7d4b796aa550 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py @@ -206,7 +206,8 @@ def __init__(self, pipeline: Pipeline, response: AsyncHttpResponse, raw: bool=Fa self.block_size = response.block_size self._raw = raw self.content_length = int(response.internal_response.headers.get('Content-Length', 0)) - if self._raw and hasattr(self.response.internal_response.raw, 'stream'): + if self._raw and hasattr(self.response.internal_response, 'raw') \ + and hasattr(self.response.internal_response.raw, 'stream'): delattr(self.response.internal_response.raw.__class__, 'stream') def __len__(self): diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_asyncio.py b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_asyncio.py index 46678e4ff040..8d56dbd62941 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_asyncio.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_asyncio.py @@ -147,7 +147,8 @@ def __init__(self, pipeline: Pipeline, response: AsyncHttpResponse, raw: bool=Fa self.block_size = response.block_size self._raw = raw self.iter_content_func = self.response.internal_response.iter_content(self.block_size) - if self._raw and hasattr(self.response.internal_response.raw, 'stream'): + if self._raw and hasattr(self.response.internal_response, 'raw') \ + and hasattr(self.response.internal_response.raw, 'stream'): delattr(self.response.internal_response.raw.__class__, 'stream') self.content_length = int(response.headers.get('Content-Length', 0)) diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_basic.py b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_basic.py index 92886d6f7c53..971359e98f4a 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_basic.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_basic.py @@ -107,7 +107,8 @@ def __init__(self, pipeline, response, raw=False): self.block_size = response.block_size self._raw = raw self.iter_content_func = self.response.internal_response.iter_content(self.block_size) - if self._raw and hasattr(self.response.internal_response.raw, 'stream'): + if self._raw and hasattr(self.response.internal_response, 'raw') \ + and hasattr(self.response.internal_response.raw, 'stream'): delattr(self.response.internal_response.raw.__class__, 'stream') self.content_length = int(response.headers.get('Content-Length', 0)) diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_trio.py b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_trio.py index eb0f2f79a65d..662d93b300d6 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_trio.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_trio.py @@ -63,7 +63,8 @@ def __init__(self, pipeline: Pipeline, response: AsyncHttpResponse, raw: bool=Fa self.block_size = response.block_size self._raw = raw self.iter_content_func = self.response.internal_response.iter_content(self.block_size) - if self._raw and hasattr(self.response.internal_response.raw, 'stream'): + if self._raw and hasattr(self.response.internal_response, 'raw') \ + and hasattr(self.response.internal_response.raw, 'stream'): delattr(self.response.internal_response.raw.__class__, 'stream') self.content_length = int(response.headers.get('Content-Length', 0)) From a6dcefa6312a5b7d649174c435338b79aec93d4e Mon Sep 17 00:00:00 2001 From: Xiang Yan Date: Wed, 14 Apr 2021 14:19:03 -0700 Subject: [PATCH 04/36] update --- sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py | 2 +- .../azure/core/pipeline/transport/_requests_asyncio.py | 2 +- .../azure-core/azure/core/pipeline/transport/_requests_trio.py | 2 +- 3 files changed, 3 insertions(+), 3 deletions(-) diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py b/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py index 7d4b796aa550..2b5dd0891f1b 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py @@ -199,7 +199,7 @@ class AioHttpStreamDownloadGenerator(AsyncIterator): :param response: The client response object. :param raw: If returns the raw stream. """ - def __init__(self, pipeline: Pipeline, response: AsyncHttpResponse, raw: bool=False) -> None: + def __init__(self, pipeline: Pipeline, response: AsyncHttpResponse, raw: bool = False) -> None: self.pipeline = pipeline self.request = response.request self.response = response diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_asyncio.py b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_asyncio.py index 8d56dbd62941..037545ec6c3c 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_asyncio.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_asyncio.py @@ -140,7 +140,7 @@ class AsyncioStreamDownloadGenerator(AsyncIterator): :param response: The response object. :param raw: If returns the raw stream. """ - def __init__(self, pipeline: Pipeline, response: AsyncHttpResponse, raw: bool=False) -> None: + def __init__(self, pipeline: Pipeline, response: AsyncHttpResponse, raw: bool = False) -> None: self.pipeline = pipeline self.request = response.request self.response = response diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_trio.py b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_trio.py index 662d93b300d6..d5e2cb3fd4f0 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_trio.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_trio.py @@ -56,7 +56,7 @@ class TrioStreamDownloadGenerator(AsyncIterator): :param response: The response object. :param raw: If returns the raw stream. """ - def __init__(self, pipeline: Pipeline, response: AsyncHttpResponse, raw: bool=False) -> None: + def __init__(self, pipeline: Pipeline, response: AsyncHttpResponse, raw: bool = False) -> None: self.pipeline = pipeline self.request = response.request self.response = response From f2d522314227941e6ba1e13fa7704e1aa5717c1c Mon Sep 17 00:00:00 2001 From: Xiang Yan Date: Thu, 15 Apr 2021 14:42:20 -0700 Subject: [PATCH 05/36] update --- .../azure/core/pipeline/transport/_aiohttp.py | 4 +-- .../pipeline/transport/_requests_asyncio.py | 25 ++++++---------- .../pipeline/transport/_requests_basic.py | 30 ++++++++++++++++--- .../core/pipeline/transport/_requests_trio.py | 10 +++---- 4 files changed, 41 insertions(+), 28 deletions(-) diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py b/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py index 2b5dd0891f1b..00595cabbe0a 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py @@ -41,6 +41,7 @@ AsyncHttpTransport, AsyncHttpResponse, _ResponseStopIteration) +from ._requests_basic import _read_raw_stream # Matching requests, because why not? CONTENT_CHUNK_SIZE = 10 * 1024 @@ -206,9 +207,6 @@ def __init__(self, pipeline: Pipeline, response: AsyncHttpResponse, raw: bool = self.block_size = response.block_size self._raw = raw self.content_length = int(response.internal_response.headers.get('Content-Length', 0)) - if self._raw and hasattr(self.response.internal_response, 'raw') \ - and hasattr(self.response.internal_response.raw, 'stream'): - delattr(self.response.internal_response.raw.__class__, 'stream') def __len__(self): return self.content_length diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_asyncio.py b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_asyncio.py index 037545ec6c3c..5cf0ea09bfa8 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_asyncio.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_asyncio.py @@ -42,7 +42,7 @@ AsyncHttpResponse, _ResponseStopIteration, _iterate_response_content) -from ._requests_basic import RequestsTransportResponse +from ._requests_basic import RequestsTransportResponse, _read_raw_stream from ._base_requests_async import RequestsAsyncTransportBase @@ -146,10 +146,10 @@ def __init__(self, pipeline: Pipeline, response: AsyncHttpResponse, raw: bool = self.response = response self.block_size = response.block_size self._raw = raw - self.iter_content_func = self.response.internal_response.iter_content(self.block_size) - if self._raw and hasattr(self.response.internal_response, 'raw') \ - and hasattr(self.response.internal_response.raw, 'stream'): - delattr(self.response.internal_response.raw.__class__, 'stream') + if self._raw: + self.iter_content_func = _read_raw_stream(self.response.internal_response, self.block_size) + else: + self.iter_content_func = self.response.internal_response.iter_content(self.block_size) self.content_length = int(response.headers.get('Content-Length', 0)) def __len__(self): @@ -158,17 +158,10 @@ def __len__(self): async def __anext__(self): loop = _get_running_loop() try: - if self._raw: - chunk = await loop.run_in_executor( - None, - _iterate_response_content, - self.iter_content_func, - ) - else: - chunk = await loop.run_in_executor( - None, - _iterate_response_content, - self.iter_content_func, + chunk = await loop.run_in_executor( + None, + _iterate_response_content, + self.iter_content_func, ) if not chunk: raise _ResponseStopIteration() diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_basic.py b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_basic.py index 971359e98f4a..e9406d7aa9b1 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_basic.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_basic.py @@ -28,6 +28,9 @@ from typing import Iterator, Optional, Any, Union, TypeVar import urllib3 # type: ignore from urllib3.util.retry import Retry # type: ignore +from urllib3.exceptions import ( + DecodeError, ReadTimeoutError, ProtocolError, LocationParseError +) import requests from azure.core.configuration import ConnectionConfiguration @@ -48,6 +51,25 @@ _LOGGER = logging.getLogger(__name__) +def _read_raw_stream(response, chunk_size=1): + # Special case for urllib3. + if hasattr(response.raw, 'stream'): + try: + for chunk in response.raw.stream(chunk_size, decode_content=False): + yield chunk + except ProtocolError as e: + raise ChunkedEncodingError(e) + except DecodeError as e: + raise ContentDecodingError(e) + except ReadTimeoutError as e: + raise ConnectionError(e) + else: + # Standard file-like object. + while True: + chunk = response.raw.read(chunk_size) + if not chunk: + break + yield chunk class _RequestsTransportResponseBase(_HttpResponseBase): """Base class for accessing response data. @@ -106,10 +128,10 @@ def __init__(self, pipeline, response, raw=False): self.response = response self.block_size = response.block_size self._raw = raw - self.iter_content_func = self.response.internal_response.iter_content(self.block_size) - if self._raw and hasattr(self.response.internal_response, 'raw') \ - and hasattr(self.response.internal_response.raw, 'stream'): - delattr(self.response.internal_response.raw.__class__, 'stream') + if self._raw: + self.iter_content_func = _read_raw_stream(self.response.internal_response, self.block_size) + else: + self.iter_content_func = self.response.internal_response.iter_content(self.block_size) self.content_length = int(response.headers.get('Content-Length', 0)) def __len__(self): diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_trio.py b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_trio.py index d5e2cb3fd4f0..6b9e144e0893 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_trio.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_trio.py @@ -42,7 +42,7 @@ AsyncHttpResponse, _ResponseStopIteration, _iterate_response_content) -from ._requests_basic import RequestsTransportResponse +from ._requests_basic import RequestsTransportResponse, _read_raw_stream from ._base_requests_async import RequestsAsyncTransportBase @@ -62,10 +62,10 @@ def __init__(self, pipeline: Pipeline, response: AsyncHttpResponse, raw: bool = self.response = response self.block_size = response.block_size self._raw = raw - self.iter_content_func = self.response.internal_response.iter_content(self.block_size) - if self._raw and hasattr(self.response.internal_response, 'raw') \ - and hasattr(self.response.internal_response.raw, 'stream'): - delattr(self.response.internal_response.raw.__class__, 'stream') + if self._raw: + self.iter_content_func = _read_raw_stream(self.response.internal_response, self.block_size) + else: + self.iter_content_func = self.response.internal_response.iter_content(self.block_size) self.content_length = int(response.headers.get('Content-Length', 0)) def __len__(self): From 8a0f2d9f8842d3266c5a45ae03444859561c7c50 Mon Sep 17 00:00:00 2001 From: Xiang Yan Date: Thu, 15 Apr 2021 15:24:15 -0700 Subject: [PATCH 06/36] update --- .../azure/core/pipeline/transport/_requests_basic.py | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_basic.py b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_basic.py index e9406d7aa9b1..c0eeaa1e83c6 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_basic.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_basic.py @@ -29,7 +29,7 @@ import urllib3 # type: ignore from urllib3.util.retry import Retry # type: ignore from urllib3.exceptions import ( - DecodeError, ReadTimeoutError, ProtocolError, LocationParseError + DecodeError, ReadTimeoutError, ProtocolError ) import requests @@ -58,11 +58,11 @@ def _read_raw_stream(response, chunk_size=1): for chunk in response.raw.stream(chunk_size, decode_content=False): yield chunk except ProtocolError as e: - raise ChunkedEncodingError(e) + raise requests.exceptions.ChunkedEncodingError(e) except DecodeError as e: - raise ContentDecodingError(e) + raise requests.exceptions.ContentDecodingError(e) except ReadTimeoutError as e: - raise ConnectionError(e) + raise requests.exceptions.ConnectionError(e) else: # Standard file-like object. while True: From 50aaa4a2b8bf2b98bf5fa02ded2af9035e8deb14 Mon Sep 17 00:00:00 2001 From: Xiang Yan Date: Thu, 15 Apr 2021 15:57:15 -0700 Subject: [PATCH 07/36] update --- sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py | 1 - 1 file changed, 1 deletion(-) diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py b/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py index 00595cabbe0a..5f6e99a4d9f4 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py @@ -41,7 +41,6 @@ AsyncHttpTransport, AsyncHttpResponse, _ResponseStopIteration) -from ._requests_basic import _read_raw_stream # Matching requests, because why not? CONTENT_CHUNK_SIZE = 10 * 1024 From 598bfa92ef4fb0272fe883c4781d57f0aaa3fb5d Mon Sep 17 00:00:00 2001 From: Xiang Yan Date: Thu, 15 Apr 2021 17:10:10 -0700 Subject: [PATCH 08/36] update --- sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py | 2 +- .../azure/core/pipeline/transport/_requests_asyncio.py | 2 +- .../azure-core/azure/core/pipeline/transport/_requests_basic.py | 2 +- .../azure-core/azure/core/pipeline/transport/_requests_trio.py | 2 +- 4 files changed, 4 insertions(+), 4 deletions(-) diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py b/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py index 5f6e99a4d9f4..f8f134a39e28 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py @@ -268,7 +268,7 @@ async def load_body(self) -> None: """Load in memory the body, so it could be accessible from sync methods.""" self._body = await self.internal_response.read() - def stream_download(self, pipeline, raw=True) -> AsyncIteratorType[bytes]: + def stream_download(self, pipeline, raw=False) -> AsyncIteratorType[bytes]: """Generator for streaming response body data. :param pipeline: The pipeline object diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_asyncio.py b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_asyncio.py index 5cf0ea09bfa8..1514df5c4c2c 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_asyncio.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_asyncio.py @@ -180,6 +180,6 @@ async def __anext__(self): class AsyncioRequestsTransportResponse(AsyncHttpResponse, RequestsTransportResponse): # type: ignore """Asynchronous streaming of data from the response. """ - def stream_download(self, pipeline, raw=True) -> AsyncIteratorType[bytes]: # type: ignore + def stream_download(self, pipeline, raw=False) -> AsyncIteratorType[bytes]: # type: ignore """Generator for streaming request body data.""" return AsyncioStreamDownloadGenerator(pipeline, self, raw=raw) # type: ignore diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_basic.py b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_basic.py index c0eeaa1e83c6..56c25500b3ab 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_basic.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_basic.py @@ -161,7 +161,7 @@ def __next__(self): class RequestsTransportResponse(HttpResponse, _RequestsTransportResponseBase): """Streaming of data from the response. """ - def stream_download(self, pipeline, raw=True): + def stream_download(self, pipeline, raw=False): # type: (PipelineType, bool) -> Iterator[bytes] """Generator for streaming request body data.""" return StreamDownloadGenerator(pipeline, self, raw=raw) diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_trio.py b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_trio.py index 6b9e144e0893..268c54d0dfac 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_trio.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_trio.py @@ -99,7 +99,7 @@ async def __anext__(self): class TrioRequestsTransportResponse(AsyncHttpResponse, RequestsTransportResponse): # type: ignore """Asynchronous streaming of data from the response. """ - def stream_download(self, pipeline, raw=True) -> AsyncIteratorType[bytes]: # type: ignore + def stream_download(self, pipeline, raw=False) -> AsyncIteratorType[bytes]: # type: ignore """Generator for streaming response data. """ return TrioStreamDownloadGenerator(pipeline, self, raw=raw) From 8d228f266a7521a2d435e2289fc9be696c6b25f3 Mon Sep 17 00:00:00 2001 From: Xiang Yan Date: Thu, 22 Apr 2021 15:54:25 -0700 Subject: [PATCH 09/36] updates --- sdk/core/azure-core/CLIENT_LIBRARY_DEVELOPER.md | 2 +- .../azure/core/pipeline/transport/_aiohttp.py | 8 +++----- .../core/pipeline/transport/_requests_asyncio.py | 13 ++++--------- .../core/pipeline/transport/_requests_basic.py | 13 ++++--------- .../azure/core/pipeline/transport/_requests_trio.py | 13 ++++--------- sdk/core/azure-core/tests/test_stream_generator.py | 11 +++++++++-- 6 files changed, 25 insertions(+), 35 deletions(-) diff --git a/sdk/core/azure-core/CLIENT_LIBRARY_DEVELOPER.md b/sdk/core/azure-core/CLIENT_LIBRARY_DEVELOPER.md index 2a7f14e93942..a7062496ffae 100644 --- a/sdk/core/azure-core/CLIENT_LIBRARY_DEVELOPER.md +++ b/sdk/core/azure-core/CLIENT_LIBRARY_DEVELOPER.md @@ -269,7 +269,7 @@ class HttpResponse(object): def text(self, encoding=None): """Return the whole body as a string.""" - def stream_download(self, pipeline, raw=False): + def stream_download(self, pipeline): """Generator for streaming request body data. Should be implemented by sub-classes if streaming download is supported. diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py b/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py index f8f134a39e28..4b6950e0f544 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py @@ -197,14 +197,12 @@ class AioHttpStreamDownloadGenerator(AsyncIterator): :param pipeline: The pipeline object :param response: The client response object. - :param raw: If returns the raw stream. """ - def __init__(self, pipeline: Pipeline, response: AsyncHttpResponse, raw: bool = False) -> None: + def __init__(self, pipeline: Pipeline, response: AsyncHttpResponse) -> None: self.pipeline = pipeline self.request = response.request self.response = response self.block_size = response.block_size - self._raw = raw self.content_length = int(response.internal_response.headers.get('Content-Length', 0)) def __len__(self): @@ -268,13 +266,13 @@ async def load_body(self) -> None: """Load in memory the body, so it could be accessible from sync methods.""" self._body = await self.internal_response.read() - def stream_download(self, pipeline, raw=False) -> AsyncIteratorType[bytes]: + def stream_download(self, pipeline) -> AsyncIteratorType[bytes]: """Generator for streaming response body data. :param pipeline: The pipeline object :type pipeline: azure.core.pipeline """ - return AioHttpStreamDownloadGenerator(pipeline, self, raw=raw) + return AioHttpStreamDownloadGenerator(pipeline, self) def __getstate__(self): # Be sure body is loaded in memory, otherwise not pickable and let it throw diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_asyncio.py b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_asyncio.py index 1514df5c4c2c..22ba06a04867 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_asyncio.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_asyncio.py @@ -138,18 +138,13 @@ class AsyncioStreamDownloadGenerator(AsyncIterator): :param pipeline: The pipeline object :param response: The response object. - :param raw: If returns the raw stream. """ - def __init__(self, pipeline: Pipeline, response: AsyncHttpResponse, raw: bool = False) -> None: + def __init__(self, pipeline: Pipeline, response: AsyncHttpResponse) -> None: self.pipeline = pipeline self.request = response.request self.response = response self.block_size = response.block_size - self._raw = raw - if self._raw: - self.iter_content_func = _read_raw_stream(self.response.internal_response, self.block_size) - else: - self.iter_content_func = self.response.internal_response.iter_content(self.block_size) + self.iter_content_func = _read_raw_stream(self.response.internal_response, self.block_size) self.content_length = int(response.headers.get('Content-Length', 0)) def __len__(self): @@ -180,6 +175,6 @@ async def __anext__(self): class AsyncioRequestsTransportResponse(AsyncHttpResponse, RequestsTransportResponse): # type: ignore """Asynchronous streaming of data from the response. """ - def stream_download(self, pipeline, raw=False) -> AsyncIteratorType[bytes]: # type: ignore + def stream_download(self, pipeline) -> AsyncIteratorType[bytes]: # type: ignore """Generator for streaming request body data.""" - return AsyncioStreamDownloadGenerator(pipeline, self, raw=raw) # type: ignore + return AsyncioStreamDownloadGenerator(pipeline, self) # type: ignore diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_basic.py b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_basic.py index 56c25500b3ab..cfceea767a64 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_basic.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_basic.py @@ -120,18 +120,13 @@ class StreamDownloadGenerator(object): :param pipeline: The pipeline object :param response: The response object. - :param raw: If returns the raw stream. """ - def __init__(self, pipeline, response, raw=False): + def __init__(self, pipeline, response): self.pipeline = pipeline self.request = response.request self.response = response self.block_size = response.block_size - self._raw = raw - if self._raw: - self.iter_content_func = _read_raw_stream(self.response.internal_response, self.block_size) - else: - self.iter_content_func = self.response.internal_response.iter_content(self.block_size) + self.iter_content_func = _read_raw_stream(self.response.internal_response, self.block_size) self.content_length = int(response.headers.get('Content-Length', 0)) def __len__(self): @@ -161,10 +156,10 @@ def __next__(self): class RequestsTransportResponse(HttpResponse, _RequestsTransportResponseBase): """Streaming of data from the response. """ - def stream_download(self, pipeline, raw=False): + def stream_download(self, pipeline): # type: (PipelineType, bool) -> Iterator[bytes] """Generator for streaming request body data.""" - return StreamDownloadGenerator(pipeline, self, raw=raw) + return StreamDownloadGenerator(pipeline, self) class RequestsTransport(HttpTransport): diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_trio.py b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_trio.py index 268c54d0dfac..8241fc7399e0 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_trio.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_trio.py @@ -54,18 +54,13 @@ class TrioStreamDownloadGenerator(AsyncIterator): :param pipeline: The pipeline object :param response: The response object. - :param raw: If returns the raw stream. """ - def __init__(self, pipeline: Pipeline, response: AsyncHttpResponse, raw: bool = False) -> None: + def __init__(self, pipeline: Pipeline, response: AsyncHttpResponse) -> None: self.pipeline = pipeline self.request = response.request self.response = response self.block_size = response.block_size - self._raw = raw - if self._raw: - self.iter_content_func = _read_raw_stream(self.response.internal_response, self.block_size) - else: - self.iter_content_func = self.response.internal_response.iter_content(self.block_size) + self.iter_content_func = _read_raw_stream(self.response.internal_response, self.block_size) self.content_length = int(response.headers.get('Content-Length', 0)) def __len__(self): @@ -99,10 +94,10 @@ async def __anext__(self): class TrioRequestsTransportResponse(AsyncHttpResponse, RequestsTransportResponse): # type: ignore """Asynchronous streaming of data from the response. """ - def stream_download(self, pipeline, raw=False) -> AsyncIteratorType[bytes]: # type: ignore + def stream_download(self, pipeline) -> AsyncIteratorType[bytes]: # type: ignore """Generator for streaming response data. """ - return TrioStreamDownloadGenerator(pipeline, self, raw=raw) + return TrioStreamDownloadGenerator(pipeline, self) class TrioRequestsTransport(RequestsAsyncTransportBase): # type: ignore diff --git a/sdk/core/azure-core/tests/test_stream_generator.py b/sdk/core/azure-core/tests/test_stream_generator.py index d3834c1b5e34..5a6f9c875f60 100644 --- a/sdk/core/azure-core/tests/test_stream_generator.py +++ b/sdk/core/azure-core/tests/test_stream_generator.py @@ -43,10 +43,17 @@ def __next__(self): if self._count == 0: self._count += 1 raise requests.exceptions.ConnectionError + + def stream(self, chunk_size, decode_content=False): + if self._count == 0: + self._count += 1 + raise requests.exceptions.ConnectionError + while True: + yield b"test" class MockInternalResponse(): - def iter_content(self, block_size): - return MockTransport() + def __init__(self): + self.raw = MockTransport() def close(self): pass From 982f083fd395e5c77d94aac9ad70b388358737ed Mon Sep 17 00:00:00 2001 From: Xiang Yan Date: Thu, 22 Apr 2021 16:01:07 -0700 Subject: [PATCH 10/36] update --- sdk/core/azure-core/azure/core/pipeline/transport/_base.py | 4 ++-- .../azure-core/azure/core/pipeline/transport/_base_async.py | 2 +- .../azure/core/pipeline/transport/_requests_basic.py | 2 +- 3 files changed, 4 insertions(+), 4 deletions(-) diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_base.py b/sdk/core/azure-core/azure/core/pipeline/transport/_base.py index 571c55623b2c..6e23d58888a4 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_base.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_base.py @@ -580,8 +580,8 @@ def __repr__(self): class HttpResponse(_HttpResponseBase): # pylint: disable=abstract-method - def stream_download(self, pipeline, raw=False): - # type: (PipelineType, bool) -> Iterator[bytes] + def stream_download(self, pipeline): + # type: (PipelineType) -> Iterator[bytes] """Generator for streaming request body data. Should be implemented by sub-classes if streaming download diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_base_async.py b/sdk/core/azure-core/azure/core/pipeline/transport/_base_async.py index 6dab35f1e992..bfc51ef6109b 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_base_async.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_base_async.py @@ -124,7 +124,7 @@ class AsyncHttpResponse(_HttpResponseBase): # pylint: disable=abstract-method Allows for the asynchronous streaming of data from the response. """ - def stream_download(self, pipeline, raw=False) -> AsyncIteratorType[bytes]: + def stream_download(self, pipeline) -> AsyncIteratorType[bytes]: """Generator for streaming response body data. Should be implemented by sub-classes if streaming download diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_basic.py b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_basic.py index cfceea767a64..834964ebb82a 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_basic.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_basic.py @@ -157,7 +157,7 @@ class RequestsTransportResponse(HttpResponse, _RequestsTransportResponseBase): """Streaming of data from the response. """ def stream_download(self, pipeline): - # type: (PipelineType, bool) -> Iterator[bytes] + # type: (PipelineType) -> Iterator[bytes] """Generator for streaming request body data.""" return StreamDownloadGenerator(pipeline, self) From 0339394888004452e9684a21cd7e67aa53103fbd Mon Sep 17 00:00:00 2001 From: Xiang Yan Date: Thu, 22 Apr 2021 16:55:56 -0700 Subject: [PATCH 11/36] update --- sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py b/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py index 4b6950e0f544..80d77fe7b13a 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py @@ -89,7 +89,8 @@ async def open(self): self.session = aiohttp.ClientSession( loop=self._loop, trust_env=self._use_env_settings, - cookie_jar=jar + cookie_jar=jar, + auto_decompress=False, ) if self.session is not None: await self.session.__aenter__() From 4c39744cf4869aedef5e5713ca713db356ea436a Mon Sep 17 00:00:00 2001 From: Xiang Yan Date: Mon, 3 May 2021 16:43:34 -0700 Subject: [PATCH 12/36] update --- .../azure-core/azure/core/pipeline/_base.py | 5 +++++ .../azure/core/pipeline/_base_async.py | 5 +++++ .../azure/core/pipeline/transport/_aiohttp.py | 19 ++++++++++++++----- .../azure/core/pipeline/transport/_base.py | 4 ++-- .../core/pipeline/transport/_base_async.py | 5 +++-- .../pipeline/transport/_requests_asyncio.py | 11 +++++++---- .../pipeline/transport/_requests_basic.py | 12 ++++++++---- .../core/pipeline/transport/_requests_trio.py | 12 ++++++++---- .../test_stream_generator_async.py | 11 ++++++++++- .../azure-core/tests/test_stream_generator.py | 2 +- 10 files changed, 63 insertions(+), 23 deletions(-) diff --git a/sdk/core/azure-core/azure/core/pipeline/_base.py b/sdk/core/azure-core/azure/core/pipeline/_base.py index b67721dec26f..a573be09eece 100644 --- a/sdk/core/azure-core/azure/core/pipeline/_base.py +++ b/sdk/core/azure-core/azure/core/pipeline/_base.py @@ -147,6 +147,11 @@ def __enter__(self): def __exit__(self, *exc_details): # pylint: disable=arguments-differ self._transport.__exit__(*exc_details) + @property + def transport(self): + """Runs transport used in the pipeline.""" + return self._transport + @staticmethod def _prepare_multipart_mixed_request(request): # type: (HTTPRequestType) -> None diff --git a/sdk/core/azure-core/azure/core/pipeline/_base_async.py b/sdk/core/azure-core/azure/core/pipeline/_base_async.py index 6d9fa19c83ef..9a5099cf36fb 100644 --- a/sdk/core/azure-core/azure/core/pipeline/_base_async.py +++ b/sdk/core/azure-core/azure/core/pipeline/_base_async.py @@ -160,6 +160,11 @@ async def __aenter__(self) -> "AsyncPipeline": async def __aexit__(self, *exc_details): # pylint: disable=arguments-differ await self._transport.__aexit__(*exc_details) + @property + def transport(self): + """Runs transport used in the pipeline.""" + return self._transport + async def _prepare_multipart_mixed_request(self, request: HTTPRequestType) -> None: """Will execute the multipart policies. diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py b/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py index 80d77fe7b13a..0a36519d88ca 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py @@ -90,7 +90,6 @@ async def open(self): loop=self._loop, trust_env=self._use_env_settings, cookie_jar=jar, - auto_decompress=False, ) if self.session is not None: await self.session.__aenter__() @@ -198,29 +197,38 @@ class AioHttpStreamDownloadGenerator(AsyncIterator): :param pipeline: The pipeline object :param response: The client response object. + :param int raw: If returns the raw stream. """ - def __init__(self, pipeline: Pipeline, response: AsyncHttpResponse) -> None: + def __init__(self, pipeline: Pipeline, response: AsyncHttpResponse, raw: bool = False) -> None: self.pipeline = pipeline self.request = response.request self.response = response self.block_size = response.block_size + self._raw = raw self.content_length = int(response.internal_response.headers.get('Content-Length', 0)) def __len__(self): return self.content_length async def __anext__(self): + raw = self.pipeline.transport.session.auto_decompress + if self._raw: + self.pipeline.transport.session.auto_decompress = False try: chunk = await self.response.internal_response.content.read(self.block_size) + self.pipeline.transport.session.auto_decompress = raw if not chunk: raise _ResponseStopIteration() return chunk except _ResponseStopIteration: + self.pipeline.transport.session.auto_decompress = raw self.response.internal_response.close() raise StopAsyncIteration() except StreamConsumedError: + self.pipeline.transport.session.auto_decompress = raw raise except Exception as err: + self.pipeline.transport.session.auto_decompress = raw _LOGGER.warning("Unable to stream download: %s", err) self.response.internal_response.close() raise @@ -267,13 +275,14 @@ async def load_body(self) -> None: """Load in memory the body, so it could be accessible from sync methods.""" self._body = await self.internal_response.read() - def stream_download(self, pipeline) -> AsyncIteratorType[bytes]: + def stream_download(self, pipeline, raw=False) -> AsyncIteratorType[bytes]: """Generator for streaming response body data. :param pipeline: The pipeline object - :type pipeline: azure.core.pipeline + :type pipeline: azure.core.pipeline.Pipeline + :param int raw: If returns the raw stream. """ - return AioHttpStreamDownloadGenerator(pipeline, self) + return AioHttpStreamDownloadGenerator(pipeline, self, raw=raw) def __getstate__(self): # Be sure body is loaded in memory, otherwise not pickable and let it throw diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_base.py b/sdk/core/azure-core/azure/core/pipeline/transport/_base.py index 6e23d58888a4..571c55623b2c 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_base.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_base.py @@ -580,8 +580,8 @@ def __repr__(self): class HttpResponse(_HttpResponseBase): # pylint: disable=abstract-method - def stream_download(self, pipeline): - # type: (PipelineType) -> Iterator[bytes] + def stream_download(self, pipeline, raw=False): + # type: (PipelineType, bool) -> Iterator[bytes] """Generator for streaming request body data. Should be implemented by sub-classes if streaming download diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_base_async.py b/sdk/core/azure-core/azure/core/pipeline/transport/_base_async.py index bfc51ef6109b..3ec8023883e6 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_base_async.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_base_async.py @@ -124,14 +124,15 @@ class AsyncHttpResponse(_HttpResponseBase): # pylint: disable=abstract-method Allows for the asynchronous streaming of data from the response. """ - def stream_download(self, pipeline) -> AsyncIteratorType[bytes]: + def stream_download(self, pipeline, raw=False) -> AsyncIteratorType[bytes]: """Generator for streaming response body data. Should be implemented by sub-classes if streaming download is supported. Will return an asynchronous generator. :param pipeline: The pipeline object - :type pipeline: azure.core.pipeline + :type pipeline: azure.core.pipeline.Pipeline + :param int raw: If returns the raw stream. """ def parts(self) -> AsyncIterator: diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_asyncio.py b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_asyncio.py index 22ba06a04867..4c669e8271c9 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_asyncio.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_asyncio.py @@ -139,12 +139,15 @@ class AsyncioStreamDownloadGenerator(AsyncIterator): :param pipeline: The pipeline object :param response: The response object. """ - def __init__(self, pipeline: Pipeline, response: AsyncHttpResponse) -> None: + def __init__(self, pipeline: Pipeline, response: AsyncHttpResponse, raw: bool = False) -> None: self.pipeline = pipeline self.request = response.request self.response = response self.block_size = response.block_size - self.iter_content_func = _read_raw_stream(self.response.internal_response, self.block_size) + if raw: + self.iter_content_func = _read_raw_stream(self.response.internal_response, self.block_size) + else: + self.iter_content_func = self.response.internal_response.iter_content(self.block_size) self.content_length = int(response.headers.get('Content-Length', 0)) def __len__(self): @@ -175,6 +178,6 @@ async def __anext__(self): class AsyncioRequestsTransportResponse(AsyncHttpResponse, RequestsTransportResponse): # type: ignore """Asynchronous streaming of data from the response. """ - def stream_download(self, pipeline) -> AsyncIteratorType[bytes]: # type: ignore + def stream_download(self, pipeline, raw=False) -> AsyncIteratorType[bytes]: # type: ignore """Generator for streaming request body data.""" - return AsyncioStreamDownloadGenerator(pipeline, self) # type: ignore + return AsyncioStreamDownloadGenerator(pipeline, self, raw=raw) # type: ignore diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_basic.py b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_basic.py index 834964ebb82a..c011417c16d6 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_basic.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_basic.py @@ -120,13 +120,17 @@ class StreamDownloadGenerator(object): :param pipeline: The pipeline object :param response: The response object. + :param int raw: If returns the raw stream. """ - def __init__(self, pipeline, response): + def __init__(self, pipeline, response, raw=False): self.pipeline = pipeline self.request = response.request self.response = response self.block_size = response.block_size - self.iter_content_func = _read_raw_stream(self.response.internal_response, self.block_size) + if raw: + self.iter_content_func = _read_raw_stream(self.response.internal_response, self.block_size) + else: + self.iter_content_func = self.response.internal_response.iter_content(self.block_size) self.content_length = int(response.headers.get('Content-Length', 0)) def __len__(self): @@ -156,10 +160,10 @@ def __next__(self): class RequestsTransportResponse(HttpResponse, _RequestsTransportResponseBase): """Streaming of data from the response. """ - def stream_download(self, pipeline): + def stream_download(self, pipeline, raw=False): # type: (PipelineType) -> Iterator[bytes] """Generator for streaming request body data.""" - return StreamDownloadGenerator(pipeline, self) + return StreamDownloadGenerator(pipeline, self, raw=raw) class RequestsTransport(HttpTransport): diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_trio.py b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_trio.py index 8241fc7399e0..b75ace652f70 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_trio.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_trio.py @@ -54,13 +54,17 @@ class TrioStreamDownloadGenerator(AsyncIterator): :param pipeline: The pipeline object :param response: The response object. + :param raw: If returns the raw stream. """ - def __init__(self, pipeline: Pipeline, response: AsyncHttpResponse) -> None: + def __init__(self, pipeline: Pipeline, response: AsyncHttpResponse, raw: bool = False) -> None: self.pipeline = pipeline self.request = response.request self.response = response self.block_size = response.block_size - self.iter_content_func = _read_raw_stream(self.response.internal_response, self.block_size) + if raw: + self.iter_content_func = _read_raw_stream(self.response.internal_response, self.block_size) + else: + self.iter_content_func = self.response.internal_response.iter_content(self.block_size) self.content_length = int(response.headers.get('Content-Length', 0)) def __len__(self): @@ -94,10 +98,10 @@ async def __anext__(self): class TrioRequestsTransportResponse(AsyncHttpResponse, RequestsTransportResponse): # type: ignore """Asynchronous streaming of data from the response. """ - def stream_download(self, pipeline) -> AsyncIteratorType[bytes]: # type: ignore + def stream_download(self, pipeline, raw=False) -> AsyncIteratorType[bytes]: # type: ignore """Generator for streaming response data. """ - return TrioStreamDownloadGenerator(pipeline, self) + return TrioStreamDownloadGenerator(pipeline, self, raw=raw) class TrioRequestsTransport(RequestsAsyncTransportBase): # type: ignore diff --git a/sdk/core/azure-core/tests/async_tests/test_stream_generator_async.py b/sdk/core/azure-core/tests/async_tests/test_stream_generator_async.py index af53b36a0045..39ecc3b065ea 100644 --- a/sdk/core/azure-core/tests/async_tests/test_stream_generator_async.py +++ b/sdk/core/azure-core/tests/async_tests/test_stream_generator_async.py @@ -17,9 +17,18 @@ @pytest.mark.asyncio async def test_connection_error_response(): + class MockSession(object): + def __init__(self): + self.auto_decompress = True + + @property + def auto_decompress(self): + return self.auto_decompress + class MockTransport(AsyncHttpTransport): def __init__(self): self._count = 0 + self.session = MockSession async def __aexit__(self, exc_type, exc_val, exc_tb): pass @@ -60,7 +69,7 @@ async def __call__(self, *args, **kwargs): pipeline = AsyncPipeline(MockTransport()) http_response = AsyncHttpResponse(http_request, None) http_response.internal_response = MockInternalResponse() - stream = AioHttpStreamDownloadGenerator(pipeline, http_response) + stream = AioHttpStreamDownloadGenerator(pipeline, http_response, raw=True) with mock.patch('asyncio.sleep', new_callable=AsyncMock): with pytest.raises(ConnectionError): await stream.__anext__() diff --git a/sdk/core/azure-core/tests/test_stream_generator.py b/sdk/core/azure-core/tests/test_stream_generator.py index 5a6f9c875f60..994f84e3c5bc 100644 --- a/sdk/core/azure-core/tests/test_stream_generator.py +++ b/sdk/core/azure-core/tests/test_stream_generator.py @@ -62,7 +62,7 @@ def close(self): pipeline = Pipeline(MockTransport()) http_response = HttpResponse(http_request, None) http_response.internal_response = MockInternalResponse() - stream = StreamDownloadGenerator(pipeline, http_response) + stream = StreamDownloadGenerator(pipeline, http_response, raw=True) with mock.patch('time.sleep', return_value=None): with pytest.raises(requests.exceptions.ConnectionError): stream.__next__() From dfcbe0ec86fd09343e1cbf0ed075c367570cc815 Mon Sep 17 00:00:00 2001 From: Xiang Yan Date: Tue, 4 May 2021 08:56:39 -0700 Subject: [PATCH 13/36] update --- .../azure-core/azure/core/pipeline/transport/_requests_basic.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_basic.py b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_basic.py index c011417c16d6..c872c16c5541 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_basic.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_basic.py @@ -161,7 +161,7 @@ class RequestsTransportResponse(HttpResponse, _RequestsTransportResponseBase): """Streaming of data from the response. """ def stream_download(self, pipeline, raw=False): - # type: (PipelineType) -> Iterator[bytes] + # type: (PipelineType, bool) -> Iterator[bytes] """Generator for streaming request body data.""" return StreamDownloadGenerator(pipeline, self, raw=raw) From 0002513e2ee70ee3193b9b8131e720b142cbd237 Mon Sep 17 00:00:00 2001 From: Xiang Yan Date: Tue, 4 May 2021 09:14:46 -0700 Subject: [PATCH 14/36] update --- .../azure/core/pipeline/transport/_aiohttp.py | 15 ++++++++++----- 1 file changed, 10 insertions(+), 5 deletions(-) diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py b/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py index 0a36519d88ca..18a7792cbce4 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py @@ -211,24 +211,29 @@ def __len__(self): return self.content_length async def __anext__(self): - raw = self.pipeline.transport.session.auto_decompress + if self.pipeline: + raw = self.pipeline.transport.session.auto_decompress if self._raw: self.pipeline.transport.session.auto_decompress = False try: chunk = await self.response.internal_response.content.read(self.block_size) - self.pipeline.transport.session.auto_decompress = raw + if self.pipeline: + self.pipeline.transport.session.auto_decompress = raw if not chunk: raise _ResponseStopIteration() return chunk except _ResponseStopIteration: - self.pipeline.transport.session.auto_decompress = raw + if self.pipeline: + self.pipeline.transport.session.auto_decompress = raw self.response.internal_response.close() raise StopAsyncIteration() except StreamConsumedError: - self.pipeline.transport.session.auto_decompress = raw + if self.pipeline: + self.pipeline.transport.session.auto_decompress = raw raise except Exception as err: - self.pipeline.transport.session.auto_decompress = raw + if self.pipeline: + self.pipeline.transport.session.auto_decompress = raw _LOGGER.warning("Unable to stream download: %s", err) self.response.internal_response.close() raise From 7d2091902c93e3664ffb58b2fb005af3dba6ddfa Mon Sep 17 00:00:00 2001 From: Xiang Yan Date: Tue, 4 May 2021 09:25:40 -0700 Subject: [PATCH 15/36] update --- .../azure/core/pipeline/transport/_aiohttp.py | 20 +++++++++++++++---- 1 file changed, 16 insertions(+), 4 deletions(-) diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py b/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py index 18a7792cbce4..a7a57d547a47 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py @@ -214,7 +214,10 @@ async def __anext__(self): if self.pipeline: raw = self.pipeline.transport.session.auto_decompress if self._raw: - self.pipeline.transport.session.auto_decompress = False + try: + self.pipeline.transport.session.auto_decompress = False + except AttributeError: + pass try: chunk = await self.response.internal_response.content.read(self.block_size) if self.pipeline: @@ -224,16 +227,25 @@ async def __anext__(self): return chunk except _ResponseStopIteration: if self.pipeline: - self.pipeline.transport.session.auto_decompress = raw + try: + self.pipeline.transport.session.auto_decompress = raw + except AttributeError: + pass self.response.internal_response.close() raise StopAsyncIteration() except StreamConsumedError: if self.pipeline: - self.pipeline.transport.session.auto_decompress = raw + try: + self.pipeline.transport.session.auto_decompress = raw + except AttributeError: + pass raise except Exception as err: if self.pipeline: - self.pipeline.transport.session.auto_decompress = raw + try: + self.pipeline.transport.session.auto_decompress = raw + except AttributeError: + pass _LOGGER.warning("Unable to stream download: %s", err) self.response.internal_response.close() raise From b2c03a328f9e49f7a89ba9a8fd9c4dabd774b0f8 Mon Sep 17 00:00:00 2001 From: Xiang Yan Date: Tue, 4 May 2021 09:37:40 -0700 Subject: [PATCH 16/36] update --- .../azure/core/pipeline/transport/_aiohttp.py | 11 +++++------ 1 file changed, 5 insertions(+), 6 deletions(-) diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py b/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py index a7a57d547a47..a2227d71ca6c 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py @@ -211,10 +211,9 @@ def __len__(self): return self.content_length async def __anext__(self): - if self.pipeline: - raw = self.pipeline.transport.session.auto_decompress - if self._raw: + if self._raw and self.pipeline: try: + raw = self.pipeline.transport.session.auto_decompress self.pipeline.transport.session.auto_decompress = False except AttributeError: pass @@ -226,7 +225,7 @@ async def __anext__(self): raise _ResponseStopIteration() return chunk except _ResponseStopIteration: - if self.pipeline: + if self._raw and self.pipeline: try: self.pipeline.transport.session.auto_decompress = raw except AttributeError: @@ -234,14 +233,14 @@ async def __anext__(self): self.response.internal_response.close() raise StopAsyncIteration() except StreamConsumedError: - if self.pipeline: + if self._raw and self.pipeline: try: self.pipeline.transport.session.auto_decompress = raw except AttributeError: pass raise except Exception as err: - if self.pipeline: + if self._raw and self.pipeline: try: self.pipeline.transport.session.auto_decompress = raw except AttributeError: From f3ef5b83003df5a392b4fb1f60b00a5e575de69f Mon Sep 17 00:00:00 2001 From: Xiang Yan Date: Tue, 4 May 2021 09:48:53 -0700 Subject: [PATCH 17/36] update --- sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py b/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py index a2227d71ca6c..4e1927c3beeb 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py @@ -219,7 +219,7 @@ async def __anext__(self): pass try: chunk = await self.response.internal_response.content.read(self.block_size) - if self.pipeline: + if self._raw and self.pipeline: self.pipeline.transport.session.auto_decompress = raw if not chunk: raise _ResponseStopIteration() From ffe9a8cfdd067bd248f1d564d48eac6fe98e91ac Mon Sep 17 00:00:00 2001 From: Xiang Yan Date: Tue, 4 May 2021 15:56:19 -0700 Subject: [PATCH 18/36] update --- .../azure-core/CLIENT_LIBRARY_DEVELOPER.md | 2 +- .../azure/core/pipeline/transport/_aiohttp.py | 34 ++++++++++--------- .../azure/core/pipeline/transport/_base.py | 2 +- .../core/pipeline/transport/_base_async.py | 5 +-- .../pipeline/transport/_requests_asyncio.py | 14 ++++---- .../pipeline/transport/_requests_basic.py | 15 ++++---- .../core/pipeline/transport/_requests_trio.py | 15 ++++---- .../test_stream_generator_async.py | 2 +- .../azure-core/tests/test_stream_generator.py | 2 +- 9 files changed, 49 insertions(+), 42 deletions(-) diff --git a/sdk/core/azure-core/CLIENT_LIBRARY_DEVELOPER.md b/sdk/core/azure-core/CLIENT_LIBRARY_DEVELOPER.md index a7062496ffae..9949d0d1ca46 100644 --- a/sdk/core/azure-core/CLIENT_LIBRARY_DEVELOPER.md +++ b/sdk/core/azure-core/CLIENT_LIBRARY_DEVELOPER.md @@ -269,7 +269,7 @@ class HttpResponse(object): def text(self, encoding=None): """Return the whole body as a string.""" - def stream_download(self, pipeline): + def stream_download(self, pipeline, decode_content=True): """Generator for streaming request body data. Should be implemented by sub-classes if streaming download is supported. diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py b/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py index 4e1927c3beeb..2a04ad64af7f 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py @@ -197,52 +197,53 @@ class AioHttpStreamDownloadGenerator(AsyncIterator): :param pipeline: The pipeline object :param response: The client response object. - :param int raw: If returns the raw stream. + :param bool decode_content: If True which is default, will attempt to decode the body based + on the ‘content-encoding’ header. """ - def __init__(self, pipeline: Pipeline, response: AsyncHttpResponse, raw: bool = False) -> None: + def __init__(self, pipeline: Pipeline, response: AsyncHttpResponse, decode_content: bool = True) -> None: self.pipeline = pipeline self.request = response.request self.response = response self.block_size = response.block_size - self._raw = raw + self._decode_content = decode_content self.content_length = int(response.internal_response.headers.get('Content-Length', 0)) def __len__(self): return self.content_length async def __anext__(self): - if self._raw and self.pipeline: + if not(self._decode_content) and self.pipeline: try: - raw = self.pipeline.transport.session.auto_decompress + auto_decompress = self.pipeline.transport.session.auto_decompress self.pipeline.transport.session.auto_decompress = False except AttributeError: pass try: chunk = await self.response.internal_response.content.read(self.block_size) - if self._raw and self.pipeline: - self.pipeline.transport.session.auto_decompress = raw + if not(self._decode_content) and self.pipeline: + self.pipeline.transport.session.auto_decompress = auto_decompress if not chunk: raise _ResponseStopIteration() return chunk except _ResponseStopIteration: - if self._raw and self.pipeline: + if not(self._decode_content) and self.pipeline: try: - self.pipeline.transport.session.auto_decompress = raw + self.pipeline.transport.session.auto_decompress = auto_decompress except AttributeError: pass self.response.internal_response.close() raise StopAsyncIteration() except StreamConsumedError: - if self._raw and self.pipeline: + if not(self._decode_content) and self.pipeline: try: - self.pipeline.transport.session.auto_decompress = raw + self.pipeline.transport.session.auto_decompress = auto_decompress except AttributeError: pass raise except Exception as err: - if self._raw and self.pipeline: + if not(self._decode_content) and self.pipeline: try: - self.pipeline.transport.session.auto_decompress = raw + self.pipeline.transport.session.auto_decompress = auto_decompress except AttributeError: pass _LOGGER.warning("Unable to stream download: %s", err) @@ -291,14 +292,15 @@ async def load_body(self) -> None: """Load in memory the body, so it could be accessible from sync methods.""" self._body = await self.internal_response.read() - def stream_download(self, pipeline, raw=False) -> AsyncIteratorType[bytes]: + def stream_download(self, pipeline, decode_content=True) -> AsyncIteratorType[bytes]: """Generator for streaming response body data. :param pipeline: The pipeline object :type pipeline: azure.core.pipeline.Pipeline - :param int raw: If returns the raw stream. + :param bool decode_content: If True which is default, will attempt to decode the body based + on the ‘content-encoding’ header. """ - return AioHttpStreamDownloadGenerator(pipeline, self, raw=raw) + return AioHttpStreamDownloadGenerator(pipeline, self, decode_content=decode_content) def __getstate__(self): # Be sure body is loaded in memory, otherwise not pickable and let it throw diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_base.py b/sdk/core/azure-core/azure/core/pipeline/transport/_base.py index 571c55623b2c..683bbdf53635 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_base.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_base.py @@ -580,7 +580,7 @@ def __repr__(self): class HttpResponse(_HttpResponseBase): # pylint: disable=abstract-method - def stream_download(self, pipeline, raw=False): + def stream_download(self, pipeline, decode_content=True): # type: (PipelineType, bool) -> Iterator[bytes] """Generator for streaming request body data. diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_base_async.py b/sdk/core/azure-core/azure/core/pipeline/transport/_base_async.py index 3ec8023883e6..ea5592090c53 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_base_async.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_base_async.py @@ -124,7 +124,7 @@ class AsyncHttpResponse(_HttpResponseBase): # pylint: disable=abstract-method Allows for the asynchronous streaming of data from the response. """ - def stream_download(self, pipeline, raw=False) -> AsyncIteratorType[bytes]: + def stream_download(self, pipeline, decode_content=True) -> AsyncIteratorType[bytes]: """Generator for streaming response body data. Should be implemented by sub-classes if streaming download @@ -132,7 +132,8 @@ def stream_download(self, pipeline, raw=False) -> AsyncIteratorType[bytes]: :param pipeline: The pipeline object :type pipeline: azure.core.pipeline.Pipeline - :param int raw: If returns the raw stream. + :param bool decode_content: If True which is default, will attempt to decode the body based + on the ‘content-encoding’ header. """ def parts(self) -> AsyncIterator: diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_asyncio.py b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_asyncio.py index 4c669e8271c9..fc51d7156704 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_asyncio.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_asyncio.py @@ -138,16 +138,18 @@ class AsyncioStreamDownloadGenerator(AsyncIterator): :param pipeline: The pipeline object :param response: The response object. + :param bool decode_content: If True which is default, will attempt to decode the body based + on the ‘content-encoding’ header. """ - def __init__(self, pipeline: Pipeline, response: AsyncHttpResponse, raw: bool = False) -> None: + def __init__(self, pipeline: Pipeline, response: AsyncHttpResponse, decode_content: bool = True) -> None: self.pipeline = pipeline self.request = response.request self.response = response self.block_size = response.block_size - if raw: - self.iter_content_func = _read_raw_stream(self.response.internal_response, self.block_size) - else: + if decode_content: self.iter_content_func = self.response.internal_response.iter_content(self.block_size) + else: + self.iter_content_func = _read_raw_stream(self.response.internal_response, self.block_size) self.content_length = int(response.headers.get('Content-Length', 0)) def __len__(self): @@ -178,6 +180,6 @@ async def __anext__(self): class AsyncioRequestsTransportResponse(AsyncHttpResponse, RequestsTransportResponse): # type: ignore """Asynchronous streaming of data from the response. """ - def stream_download(self, pipeline, raw=False) -> AsyncIteratorType[bytes]: # type: ignore + def stream_download(self, pipeline, decode_content=True) -> AsyncIteratorType[bytes]: # type: ignore """Generator for streaming request body data.""" - return AsyncioStreamDownloadGenerator(pipeline, self, raw=raw) # type: ignore + return AsyncioStreamDownloadGenerator(pipeline, self, decode_content=decode_content) # type: ignore diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_basic.py b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_basic.py index c872c16c5541..9e7270996b62 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_basic.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_basic.py @@ -120,17 +120,18 @@ class StreamDownloadGenerator(object): :param pipeline: The pipeline object :param response: The response object. - :param int raw: If returns the raw stream. + :param bool decode_content: If True which is default, will attempt to decode the body based + on the ‘content-encoding’ header. """ - def __init__(self, pipeline, response, raw=False): + def __init__(self, pipeline, response, decode_content=True): self.pipeline = pipeline self.request = response.request self.response = response self.block_size = response.block_size - if raw: - self.iter_content_func = _read_raw_stream(self.response.internal_response, self.block_size) - else: + if decode_content: self.iter_content_func = self.response.internal_response.iter_content(self.block_size) + else: + self.iter_content_func = _read_raw_stream(self.response.internal_response, self.block_size) self.content_length = int(response.headers.get('Content-Length', 0)) def __len__(self): @@ -160,10 +161,10 @@ def __next__(self): class RequestsTransportResponse(HttpResponse, _RequestsTransportResponseBase): """Streaming of data from the response. """ - def stream_download(self, pipeline, raw=False): + def stream_download(self, pipeline, decode_content=True): # type: (PipelineType, bool) -> Iterator[bytes] """Generator for streaming request body data.""" - return StreamDownloadGenerator(pipeline, self, raw=raw) + return StreamDownloadGenerator(pipeline, self, decode_content=decode_content) class RequestsTransport(HttpTransport): diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_trio.py b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_trio.py index b75ace652f70..62898867aa0e 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_trio.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_trio.py @@ -54,17 +54,18 @@ class TrioStreamDownloadGenerator(AsyncIterator): :param pipeline: The pipeline object :param response: The response object. - :param raw: If returns the raw stream. + :param bool decode_content: If True which is default, will attempt to decode the body based + on the ‘content-encoding’ header. """ - def __init__(self, pipeline: Pipeline, response: AsyncHttpResponse, raw: bool = False) -> None: + def __init__(self, pipeline: Pipeline, response: AsyncHttpResponse, decode_content: bool = True) -> None: self.pipeline = pipeline self.request = response.request self.response = response self.block_size = response.block_size - if raw: - self.iter_content_func = _read_raw_stream(self.response.internal_response, self.block_size) - else: + if decode_content: self.iter_content_func = self.response.internal_response.iter_content(self.block_size) + else: + self.iter_content_func = _read_raw_stream(self.response.internal_response, self.block_size) self.content_length = int(response.headers.get('Content-Length', 0)) def __len__(self): @@ -98,10 +99,10 @@ async def __anext__(self): class TrioRequestsTransportResponse(AsyncHttpResponse, RequestsTransportResponse): # type: ignore """Asynchronous streaming of data from the response. """ - def stream_download(self, pipeline, raw=False) -> AsyncIteratorType[bytes]: # type: ignore + def stream_download(self, pipeline, decode_content=True) -> AsyncIteratorType[bytes]: # type: ignore """Generator for streaming response data. """ - return TrioStreamDownloadGenerator(pipeline, self, raw=raw) + return TrioStreamDownloadGenerator(pipeline, self, decode_content=decode_content) class TrioRequestsTransport(RequestsAsyncTransportBase): # type: ignore diff --git a/sdk/core/azure-core/tests/async_tests/test_stream_generator_async.py b/sdk/core/azure-core/tests/async_tests/test_stream_generator_async.py index 39ecc3b065ea..3f28ab1dc2ba 100644 --- a/sdk/core/azure-core/tests/async_tests/test_stream_generator_async.py +++ b/sdk/core/azure-core/tests/async_tests/test_stream_generator_async.py @@ -69,7 +69,7 @@ async def __call__(self, *args, **kwargs): pipeline = AsyncPipeline(MockTransport()) http_response = AsyncHttpResponse(http_request, None) http_response.internal_response = MockInternalResponse() - stream = AioHttpStreamDownloadGenerator(pipeline, http_response, raw=True) + stream = AioHttpStreamDownloadGenerator(pipeline, http_response, decode_content=False) with mock.patch('asyncio.sleep', new_callable=AsyncMock): with pytest.raises(ConnectionError): await stream.__anext__() diff --git a/sdk/core/azure-core/tests/test_stream_generator.py b/sdk/core/azure-core/tests/test_stream_generator.py index 994f84e3c5bc..c91cd53047c8 100644 --- a/sdk/core/azure-core/tests/test_stream_generator.py +++ b/sdk/core/azure-core/tests/test_stream_generator.py @@ -62,7 +62,7 @@ def close(self): pipeline = Pipeline(MockTransport()) http_response = HttpResponse(http_request, None) http_response.internal_response = MockInternalResponse() - stream = StreamDownloadGenerator(pipeline, http_response, raw=True) + stream = StreamDownloadGenerator(pipeline, http_response, decode_content=False) with mock.patch('time.sleep', return_value=None): with pytest.raises(requests.exceptions.ConnectionError): stream.__next__() From 5aa8c045fda5456eb726112adf362e01d595ef5b Mon Sep 17 00:00:00 2001 From: Xiang Yan Date: Tue, 4 May 2021 16:35:28 -0700 Subject: [PATCH 19/36] update --- .../azure-core/CLIENT_LIBRARY_DEVELOPER.md | 2 +- .../azure/core/pipeline/transport/_aiohttp.py | 22 +++++++++---------- .../azure/core/pipeline/transport/_base.py | 2 +- .../pipeline/transport/_requests_asyncio.py | 10 ++++----- .../pipeline/transport/_requests_basic.py | 12 +++++----- .../core/pipeline/transport/_requests_trio.py | 10 ++++----- .../test_stream_generator_async.py | 6 ++--- .../azure-core/tests/test_stream_generator.py | 8 +++---- 8 files changed, 36 insertions(+), 36 deletions(-) diff --git a/sdk/core/azure-core/CLIENT_LIBRARY_DEVELOPER.md b/sdk/core/azure-core/CLIENT_LIBRARY_DEVELOPER.md index 9949d0d1ca46..7949e02f1627 100644 --- a/sdk/core/azure-core/CLIENT_LIBRARY_DEVELOPER.md +++ b/sdk/core/azure-core/CLIENT_LIBRARY_DEVELOPER.md @@ -269,7 +269,7 @@ class HttpResponse(object): def text(self, encoding=None): """Return the whole body as a string.""" - def stream_download(self, pipeline, decode_content=True): + def stream_download(self, pipeline, decompress=True): """Generator for streaming request body data. Should be implemented by sub-classes if streaming download is supported. diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py b/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py index 2a04ad64af7f..a49a6ff1324c 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py @@ -197,22 +197,22 @@ class AioHttpStreamDownloadGenerator(AsyncIterator): :param pipeline: The pipeline object :param response: The client response object. - :param bool decode_content: If True which is default, will attempt to decode the body based + :param bool decompress: If True which is default, will attempt to decode the body based on the ‘content-encoding’ header. """ - def __init__(self, pipeline: Pipeline, response: AsyncHttpResponse, decode_content: bool = True) -> None: + def __init__(self, pipeline: Pipeline, response: AsyncHttpResponse, decompress: bool = True) -> None: self.pipeline = pipeline self.request = response.request self.response = response self.block_size = response.block_size - self._decode_content = decode_content + self._decompress = decompress self.content_length = int(response.internal_response.headers.get('Content-Length', 0)) def __len__(self): return self.content_length async def __anext__(self): - if not(self._decode_content) and self.pipeline: + if not(self._decompress) and self.pipeline: try: auto_decompress = self.pipeline.transport.session.auto_decompress self.pipeline.transport.session.auto_decompress = False @@ -220,13 +220,13 @@ async def __anext__(self): pass try: chunk = await self.response.internal_response.content.read(self.block_size) - if not(self._decode_content) and self.pipeline: + if not(self._decompress) and self.pipeline: self.pipeline.transport.session.auto_decompress = auto_decompress if not chunk: raise _ResponseStopIteration() return chunk except _ResponseStopIteration: - if not(self._decode_content) and self.pipeline: + if not(self._decompress) and self.pipeline: try: self.pipeline.transport.session.auto_decompress = auto_decompress except AttributeError: @@ -234,14 +234,14 @@ async def __anext__(self): self.response.internal_response.close() raise StopAsyncIteration() except StreamConsumedError: - if not(self._decode_content) and self.pipeline: + if not(self._decompress) and self.pipeline: try: self.pipeline.transport.session.auto_decompress = auto_decompress except AttributeError: pass raise except Exception as err: - if not(self._decode_content) and self.pipeline: + if not(self._decompress) and self.pipeline: try: self.pipeline.transport.session.auto_decompress = auto_decompress except AttributeError: @@ -292,15 +292,15 @@ async def load_body(self) -> None: """Load in memory the body, so it could be accessible from sync methods.""" self._body = await self.internal_response.read() - def stream_download(self, pipeline, decode_content=True) -> AsyncIteratorType[bytes]: + def stream_download(self, pipeline, decompress=True) -> AsyncIteratorType[bytes]: """Generator for streaming response body data. :param pipeline: The pipeline object :type pipeline: azure.core.pipeline.Pipeline - :param bool decode_content: If True which is default, will attempt to decode the body based + :param bool decompress: If True which is default, will attempt to decode the body based on the ‘content-encoding’ header. """ - return AioHttpStreamDownloadGenerator(pipeline, self, decode_content=decode_content) + return AioHttpStreamDownloadGenerator(pipeline, self, decompress=decompress) def __getstate__(self): # Be sure body is loaded in memory, otherwise not pickable and let it throw diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_base.py b/sdk/core/azure-core/azure/core/pipeline/transport/_base.py index 683bbdf53635..6c2a71ad5280 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_base.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_base.py @@ -580,7 +580,7 @@ def __repr__(self): class HttpResponse(_HttpResponseBase): # pylint: disable=abstract-method - def stream_download(self, pipeline, decode_content=True): + def stream_download(self, pipeline, decompress=True): # type: (PipelineType, bool) -> Iterator[bytes] """Generator for streaming request body data. diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_asyncio.py b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_asyncio.py index fc51d7156704..27e2a3d62031 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_asyncio.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_asyncio.py @@ -138,15 +138,15 @@ class AsyncioStreamDownloadGenerator(AsyncIterator): :param pipeline: The pipeline object :param response: The response object. - :param bool decode_content: If True which is default, will attempt to decode the body based + :param bool decompress: If True which is default, will attempt to decode the body based on the ‘content-encoding’ header. """ - def __init__(self, pipeline: Pipeline, response: AsyncHttpResponse, decode_content: bool = True) -> None: + def __init__(self, pipeline: Pipeline, response: AsyncHttpResponse, decompress: bool = True) -> None: self.pipeline = pipeline self.request = response.request self.response = response self.block_size = response.block_size - if decode_content: + if decompress: self.iter_content_func = self.response.internal_response.iter_content(self.block_size) else: self.iter_content_func = _read_raw_stream(self.response.internal_response, self.block_size) @@ -180,6 +180,6 @@ async def __anext__(self): class AsyncioRequestsTransportResponse(AsyncHttpResponse, RequestsTransportResponse): # type: ignore """Asynchronous streaming of data from the response. """ - def stream_download(self, pipeline, decode_content=True) -> AsyncIteratorType[bytes]: # type: ignore + def stream_download(self, pipeline, decompress=True) -> AsyncIteratorType[bytes]: # type: ignore """Generator for streaming request body data.""" - return AsyncioStreamDownloadGenerator(pipeline, self, decode_content=decode_content) # type: ignore + return AsyncioStreamDownloadGenerator(pipeline, self, decompress=decompress) # type: ignore diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_basic.py b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_basic.py index 9e7270996b62..4be396003fe9 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_basic.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_basic.py @@ -55,7 +55,7 @@ def _read_raw_stream(response, chunk_size=1): # Special case for urllib3. if hasattr(response.raw, 'stream'): try: - for chunk in response.raw.stream(chunk_size, decode_content=False): + for chunk in response.raw.stream(chunk_size, decompress=False): yield chunk except ProtocolError as e: raise requests.exceptions.ChunkedEncodingError(e) @@ -120,15 +120,15 @@ class StreamDownloadGenerator(object): :param pipeline: The pipeline object :param response: The response object. - :param bool decode_content: If True which is default, will attempt to decode the body based + :param bool decompress: If True which is default, will attempt to decode the body based on the ‘content-encoding’ header. """ - def __init__(self, pipeline, response, decode_content=True): + def __init__(self, pipeline, response, decompress=True): self.pipeline = pipeline self.request = response.request self.response = response self.block_size = response.block_size - if decode_content: + if decompress: self.iter_content_func = self.response.internal_response.iter_content(self.block_size) else: self.iter_content_func = _read_raw_stream(self.response.internal_response, self.block_size) @@ -161,10 +161,10 @@ def __next__(self): class RequestsTransportResponse(HttpResponse, _RequestsTransportResponseBase): """Streaming of data from the response. """ - def stream_download(self, pipeline, decode_content=True): + def stream_download(self, pipeline, decompress=True): # type: (PipelineType, bool) -> Iterator[bytes] """Generator for streaming request body data.""" - return StreamDownloadGenerator(pipeline, self, decode_content=decode_content) + return StreamDownloadGenerator(pipeline, self, decompress=decompress) class RequestsTransport(HttpTransport): diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_trio.py b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_trio.py index 62898867aa0e..197a1529eff8 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_trio.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_trio.py @@ -54,15 +54,15 @@ class TrioStreamDownloadGenerator(AsyncIterator): :param pipeline: The pipeline object :param response: The response object. - :param bool decode_content: If True which is default, will attempt to decode the body based + :param bool decompress: If True which is default, will attempt to decode the body based on the ‘content-encoding’ header. """ - def __init__(self, pipeline: Pipeline, response: AsyncHttpResponse, decode_content: bool = True) -> None: + def __init__(self, pipeline: Pipeline, response: AsyncHttpResponse, decompress: bool = True) -> None: self.pipeline = pipeline self.request = response.request self.response = response self.block_size = response.block_size - if decode_content: + if decompress: self.iter_content_func = self.response.internal_response.iter_content(self.block_size) else: self.iter_content_func = _read_raw_stream(self.response.internal_response, self.block_size) @@ -99,10 +99,10 @@ async def __anext__(self): class TrioRequestsTransportResponse(AsyncHttpResponse, RequestsTransportResponse): # type: ignore """Asynchronous streaming of data from the response. """ - def stream_download(self, pipeline, decode_content=True) -> AsyncIteratorType[bytes]: # type: ignore + def stream_download(self, pipeline, decompress=True) -> AsyncIteratorType[bytes]: # type: ignore """Generator for streaming response data. """ - return TrioStreamDownloadGenerator(pipeline, self, decode_content=decode_content) + return TrioStreamDownloadGenerator(pipeline, self, decompress=decompress) class TrioRequestsTransport(RequestsAsyncTransportBase): # type: ignore diff --git a/sdk/core/azure-core/tests/async_tests/test_stream_generator_async.py b/sdk/core/azure-core/tests/async_tests/test_stream_generator_async.py index 3f28ab1dc2ba..77e07741e459 100644 --- a/sdk/core/azure-core/tests/async_tests/test_stream_generator_async.py +++ b/sdk/core/azure-core/tests/async_tests/test_stream_generator_async.py @@ -69,7 +69,7 @@ async def __call__(self, *args, **kwargs): pipeline = AsyncPipeline(MockTransport()) http_response = AsyncHttpResponse(http_request, None) http_response.internal_response = MockInternalResponse() - stream = AioHttpStreamDownloadGenerator(pipeline, http_response, decode_content=False) + stream = AioHttpStreamDownloadGenerator(pipeline, http_response, decompress=False) with mock.patch('asyncio.sleep', new_callable=AsyncMock): with pytest.raises(ConnectionError): await stream.__anext__() @@ -87,7 +87,7 @@ class FakeStreamWithConnectionError: def __init__(self): self.total_response_size = 500 - def stream(self, chunk_size, decode_content=False): + def stream(self, chunk_size, decompress=False): assert chunk_size == block_size left = total_response_size while left > 0: @@ -97,7 +97,7 @@ def stream(self, chunk_size, decode_content=False): left -= len(data) yield data - def read(self, chunk_size, decode_content=False): + def read(self, chunk_size, decompress=False): assert chunk_size == block_size if self.total_response_size > 0: if self.total_response_size <= block_size: diff --git a/sdk/core/azure-core/tests/test_stream_generator.py b/sdk/core/azure-core/tests/test_stream_generator.py index c91cd53047c8..2dc1b85629d9 100644 --- a/sdk/core/azure-core/tests/test_stream_generator.py +++ b/sdk/core/azure-core/tests/test_stream_generator.py @@ -44,7 +44,7 @@ def __next__(self): self._count += 1 raise requests.exceptions.ConnectionError - def stream(self, chunk_size, decode_content=False): + def stream(self, chunk_size, decompress=False): if self._count == 0: self._count += 1 raise requests.exceptions.ConnectionError @@ -62,7 +62,7 @@ def close(self): pipeline = Pipeline(MockTransport()) http_response = HttpResponse(http_request, None) http_response.internal_response = MockInternalResponse() - stream = StreamDownloadGenerator(pipeline, http_response, decode_content=False) + stream = StreamDownloadGenerator(pipeline, http_response, decompress=False) with mock.patch('time.sleep', return_value=None): with pytest.raises(requests.exceptions.ConnectionError): stream.__next__() @@ -79,7 +79,7 @@ class FakeStreamWithConnectionError: def __init__(self): self.total_response_size = 500 - def stream(self, chunk_size, decode_content=False): + def stream(self, chunk_size, decompress=False): assert chunk_size == block_size left = total_response_size while left > 0: @@ -89,7 +89,7 @@ def stream(self, chunk_size, decode_content=False): left -= len(data) yield data - def read(self, chunk_size, decode_content=False): + def read(self, chunk_size, decompress=False): assert chunk_size == block_size if self.total_response_size > 0: if self.total_response_size <= block_size: From f592f6252dc71e13e7bfdbdfca16b000d3e60c50 Mon Sep 17 00:00:00 2001 From: Xiang Yan Date: Tue, 4 May 2021 17:08:02 -0700 Subject: [PATCH 20/36] update --- .../azure-core/azure/core/pipeline/transport/_base_async.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_base_async.py b/sdk/core/azure-core/azure/core/pipeline/transport/_base_async.py index ea5592090c53..e121c95f8c68 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_base_async.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_base_async.py @@ -124,7 +124,7 @@ class AsyncHttpResponse(_HttpResponseBase): # pylint: disable=abstract-method Allows for the asynchronous streaming of data from the response. """ - def stream_download(self, pipeline, decode_content=True) -> AsyncIteratorType[bytes]: + def stream_download(self, pipeline, decompress=True) -> AsyncIteratorType[bytes]: """Generator for streaming response body data. Should be implemented by sub-classes if streaming download @@ -132,7 +132,7 @@ def stream_download(self, pipeline, decode_content=True) -> AsyncIteratorType[by :param pipeline: The pipeline object :type pipeline: azure.core.pipeline.Pipeline - :param bool decode_content: If True which is default, will attempt to decode the body based + :param bool decompress: If True which is default, will attempt to decode the body based on the ‘content-encoding’ header. """ From 68b27cf521e4c7037ad6fed47237d7b05b33d2e2 Mon Sep 17 00:00:00 2001 From: Xiang Yan Date: Wed, 5 May 2021 08:48:59 -0700 Subject: [PATCH 21/36] update --- .../azure/core/pipeline/transport/_requests_basic.py | 2 +- .../tests/async_tests/test_stream_generator_async.py | 4 ++-- sdk/core/azure-core/tests/test_stream_generator.py | 8 ++++---- 3 files changed, 7 insertions(+), 7 deletions(-) diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_basic.py b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_basic.py index 4be396003fe9..67a132926793 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_basic.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_basic.py @@ -55,7 +55,7 @@ def _read_raw_stream(response, chunk_size=1): # Special case for urllib3. if hasattr(response.raw, 'stream'): try: - for chunk in response.raw.stream(chunk_size, decompress=False): + for chunk in response.raw.stream(chunk_size, decode_content=False): yield chunk except ProtocolError as e: raise requests.exceptions.ChunkedEncodingError(e) diff --git a/sdk/core/azure-core/tests/async_tests/test_stream_generator_async.py b/sdk/core/azure-core/tests/async_tests/test_stream_generator_async.py index 77e07741e459..de7bc894e42d 100644 --- a/sdk/core/azure-core/tests/async_tests/test_stream_generator_async.py +++ b/sdk/core/azure-core/tests/async_tests/test_stream_generator_async.py @@ -87,7 +87,7 @@ class FakeStreamWithConnectionError: def __init__(self): self.total_response_size = 500 - def stream(self, chunk_size, decompress=False): + def stream(self, chunk_size, decode_content=False): assert chunk_size == block_size left = total_response_size while left > 0: @@ -97,7 +97,7 @@ def stream(self, chunk_size, decompress=False): left -= len(data) yield data - def read(self, chunk_size, decompress=False): + def read(self, chunk_size, decode_content=False): assert chunk_size == block_size if self.total_response_size > 0: if self.total_response_size <= block_size: diff --git a/sdk/core/azure-core/tests/test_stream_generator.py b/sdk/core/azure-core/tests/test_stream_generator.py index 2dc1b85629d9..c43053eeab9d 100644 --- a/sdk/core/azure-core/tests/test_stream_generator.py +++ b/sdk/core/azure-core/tests/test_stream_generator.py @@ -44,7 +44,7 @@ def __next__(self): self._count += 1 raise requests.exceptions.ConnectionError - def stream(self, chunk_size, decompress=False): + def stream(self, chunk_size, decode_content=False): if self._count == 0: self._count += 1 raise requests.exceptions.ConnectionError @@ -79,7 +79,7 @@ class FakeStreamWithConnectionError: def __init__(self): self.total_response_size = 500 - def stream(self, chunk_size, decompress=False): + def stream(self, chunk_size, decode_content=False): assert chunk_size == block_size left = total_response_size while left > 0: @@ -89,7 +89,7 @@ def stream(self, chunk_size, decompress=False): left -= len(data) yield data - def read(self, chunk_size, decompress=False): + def read(self, chunk_size, decode_content=False): assert chunk_size == block_size if self.total_response_size > 0: if self.total_response_size <= block_size: @@ -120,6 +120,6 @@ def mock_run(self, *args, **kwargs): transport = RequestsTransport() pipeline = Pipeline(transport) pipeline.run = mock_run - downloader = response.stream_download(pipeline) + downloader = response.stream_download(pipeline, decompress=False) with pytest.raises(requests.exceptions.ConnectionError): full_response = b"".join(downloader) From b6975da1607969aa8b70debca329df029a5e0265 Mon Sep 17 00:00:00 2001 From: Xiang Yan Date: Wed, 5 May 2021 11:40:33 -0700 Subject: [PATCH 22/36] update --- .../azure-core/azure/core/pipeline/_base.py | 5 -- .../azure/core/pipeline/_base_async.py | 5 -- .../azure/core/pipeline/transport/_aiohttp.py | 58 +++++++++---------- .../azure/core/pipeline/transport/_base.py | 4 +- .../core/pipeline/transport/_base_async.py | 4 +- .../pipeline/transport/_requests_asyncio.py | 9 +-- .../pipeline/transport/_requests_basic.py | 11 ++-- .../core/pipeline/transport/_requests_trio.py | 9 +-- 8 files changed, 49 insertions(+), 56 deletions(-) diff --git a/sdk/core/azure-core/azure/core/pipeline/_base.py b/sdk/core/azure-core/azure/core/pipeline/_base.py index a573be09eece..b67721dec26f 100644 --- a/sdk/core/azure-core/azure/core/pipeline/_base.py +++ b/sdk/core/azure-core/azure/core/pipeline/_base.py @@ -147,11 +147,6 @@ def __enter__(self): def __exit__(self, *exc_details): # pylint: disable=arguments-differ self._transport.__exit__(*exc_details) - @property - def transport(self): - """Runs transport used in the pipeline.""" - return self._transport - @staticmethod def _prepare_multipart_mixed_request(request): # type: (HTTPRequestType) -> None diff --git a/sdk/core/azure-core/azure/core/pipeline/_base_async.py b/sdk/core/azure-core/azure/core/pipeline/_base_async.py index 9a5099cf36fb..6d9fa19c83ef 100644 --- a/sdk/core/azure-core/azure/core/pipeline/_base_async.py +++ b/sdk/core/azure-core/azure/core/pipeline/_base_async.py @@ -160,11 +160,6 @@ async def __aenter__(self) -> "AsyncPipeline": async def __aexit__(self, *exc_details): # pylint: disable=arguments-differ await self._transport.__aexit__(*exc_details) - @property - def transport(self): - """Runs transport used in the pipeline.""" - return self._transport - async def _prepare_multipart_mixed_request(self, request: HTTPRequestType) -> None: """Will execute the multipart policies. diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py b/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py index a49a6ff1324c..f9ddc0829986 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py @@ -33,7 +33,7 @@ from requests.exceptions import StreamConsumedError from azure.core.configuration import ConnectionConfiguration -from azure.core.exceptions import ServiceRequestError, ServiceResponseError +from azure.core.exceptions import ServiceRequestError, ServiceResponseError, DecodeError from azure.core.pipeline import Pipeline from ._base import HttpRequest @@ -90,6 +90,7 @@ async def open(self): loop=self._loop, trust_env=self._use_env_settings, cookie_jar=jar, + auto_decompress=False, ) if self.session is not None: await self.session.__aenter__() @@ -197,55 +198,54 @@ class AioHttpStreamDownloadGenerator(AsyncIterator): :param pipeline: The pipeline object :param response: The client response object. - :param bool decompress: If True which is default, will attempt to decode the body based + :keyword bool decompress: If True which is default, will attempt to decode the body based on the ‘content-encoding’ header. """ - def __init__(self, pipeline: Pipeline, response: AsyncHttpResponse, decompress: bool = True) -> None: + def __init__(self, pipeline: Pipeline, response: AsyncHttpResponse, **kwargs) -> None: self.pipeline = pipeline self.request = response.request self.response = response self.block_size = response.block_size - self._decompress = decompress + self._decompress = kwargs.get("decompress", True) self.content_length = int(response.internal_response.headers.get('Content-Length', 0)) def __len__(self): return self.content_length async def __anext__(self): - if not(self._decompress) and self.pipeline: - try: - auto_decompress = self.pipeline.transport.session.auto_decompress - self.pipeline.transport.session.auto_decompress = False - except AttributeError: - pass try: chunk = await self.response.internal_response.content.read(self.block_size) - if not(self._decompress) and self.pipeline: - self.pipeline.transport.session.auto_decompress = auto_decompress if not chunk: raise _ResponseStopIteration() + if not(self._decompress): + return chunk + enc = response.internal_response.headers.get('Content-Encoding') + if not(enc): + return chunk + enc = enc.lower() + if enc in ("gzip", "deflate", "br"): + if encoding == "br": + try: + import brotli + decompressor = brotli.Decompressor() + except ImportError: + raise DecodeError( + "Can not decode content-encoding: brotli (br). " + "Please install `brotlipy`" + ) + decompressor = brotli.Decompressor() + else: + import zlib + zlib_mode = 16 + zlib.MAX_WBITS if encoding == "gzip" else zlib.MAX_WBITS + decompressor = zlib.decompressobj(wbits=zlib_mode) + chunk = decompressor.decompress(chunk) return chunk except _ResponseStopIteration: - if not(self._decompress) and self.pipeline: - try: - self.pipeline.transport.session.auto_decompress = auto_decompress - except AttributeError: - pass self.response.internal_response.close() raise StopAsyncIteration() except StreamConsumedError: - if not(self._decompress) and self.pipeline: - try: - self.pipeline.transport.session.auto_decompress = auto_decompress - except AttributeError: - pass raise except Exception as err: - if not(self._decompress) and self.pipeline: - try: - self.pipeline.transport.session.auto_decompress = auto_decompress - except AttributeError: - pass _LOGGER.warning("Unable to stream download: %s", err) self.response.internal_response.close() raise @@ -292,12 +292,12 @@ async def load_body(self) -> None: """Load in memory the body, so it could be accessible from sync methods.""" self._body = await self.internal_response.read() - def stream_download(self, pipeline, decompress=True) -> AsyncIteratorType[bytes]: + def stream_download(self, pipeline, **kwargs) -> AsyncIteratorType[bytes]: """Generator for streaming response body data. :param pipeline: The pipeline object :type pipeline: azure.core.pipeline.Pipeline - :param bool decompress: If True which is default, will attempt to decode the body based + :keyword bool decompress: If True which is default, will attempt to decode the body based on the ‘content-encoding’ header. """ return AioHttpStreamDownloadGenerator(pipeline, self, decompress=decompress) diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_base.py b/sdk/core/azure-core/azure/core/pipeline/transport/_base.py index 6c2a71ad5280..589d5549c584 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_base.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_base.py @@ -580,8 +580,8 @@ def __repr__(self): class HttpResponse(_HttpResponseBase): # pylint: disable=abstract-method - def stream_download(self, pipeline, decompress=True): - # type: (PipelineType, bool) -> Iterator[bytes] + def stream_download(self, pipeline, **kwargs): + # type: (PipelineType, **Any) -> Iterator[bytes] """Generator for streaming request body data. Should be implemented by sub-classes if streaming download diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_base_async.py b/sdk/core/azure-core/azure/core/pipeline/transport/_base_async.py index e121c95f8c68..adf09cc10a51 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_base_async.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_base_async.py @@ -124,7 +124,7 @@ class AsyncHttpResponse(_HttpResponseBase): # pylint: disable=abstract-method Allows for the asynchronous streaming of data from the response. """ - def stream_download(self, pipeline, decompress=True) -> AsyncIteratorType[bytes]: + def stream_download(self, pipeline, **kwargs) -> AsyncIteratorType[bytes]: """Generator for streaming response body data. Should be implemented by sub-classes if streaming download @@ -132,7 +132,7 @@ def stream_download(self, pipeline, decompress=True) -> AsyncIteratorType[bytes] :param pipeline: The pipeline object :type pipeline: azure.core.pipeline.Pipeline - :param bool decompress: If True which is default, will attempt to decode the body based + :keyword bool decompress: If True which is default, will attempt to decode the body based on the ‘content-encoding’ header. """ diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_asyncio.py b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_asyncio.py index 27e2a3d62031..669be477b552 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_asyncio.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_asyncio.py @@ -138,14 +138,15 @@ class AsyncioStreamDownloadGenerator(AsyncIterator): :param pipeline: The pipeline object :param response: The response object. - :param bool decompress: If True which is default, will attempt to decode the body based + :keyword bool decompress: If True which is default, will attempt to decode the body based on the ‘content-encoding’ header. """ - def __init__(self, pipeline: Pipeline, response: AsyncHttpResponse, decompress: bool = True) -> None: + def __init__(self, pipeline: Pipeline, response: AsyncHttpResponse, **kwargs) -> None: self.pipeline = pipeline self.request = response.request self.response = response self.block_size = response.block_size + decompress = kwargs.get("decompress", True) if decompress: self.iter_content_func = self.response.internal_response.iter_content(self.block_size) else: @@ -180,6 +181,6 @@ async def __anext__(self): class AsyncioRequestsTransportResponse(AsyncHttpResponse, RequestsTransportResponse): # type: ignore """Asynchronous streaming of data from the response. """ - def stream_download(self, pipeline, decompress=True) -> AsyncIteratorType[bytes]: # type: ignore + def stream_download(self, pipeline, **kwargs) -> AsyncIteratorType[bytes]: # type: ignore """Generator for streaming request body data.""" - return AsyncioStreamDownloadGenerator(pipeline, self, decompress=decompress) # type: ignore + return AsyncioStreamDownloadGenerator(pipeline, self, **kwargs) # type: ignore diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_basic.py b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_basic.py index 67a132926793..2d44ab0cd70c 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_basic.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_basic.py @@ -120,14 +120,15 @@ class StreamDownloadGenerator(object): :param pipeline: The pipeline object :param response: The response object. - :param bool decompress: If True which is default, will attempt to decode the body based + :keyword bool decompress: If True which is default, will attempt to decode the body based on the ‘content-encoding’ header. """ - def __init__(self, pipeline, response, decompress=True): + def __init__(self, pipeline, response, **kwargs): self.pipeline = pipeline self.request = response.request self.response = response self.block_size = response.block_size + decompress = kwargs.get("decompress", True) if decompress: self.iter_content_func = self.response.internal_response.iter_content(self.block_size) else: @@ -161,10 +162,10 @@ def __next__(self): class RequestsTransportResponse(HttpResponse, _RequestsTransportResponseBase): """Streaming of data from the response. """ - def stream_download(self, pipeline, decompress=True): - # type: (PipelineType, bool) -> Iterator[bytes] + def stream_download(self, pipeline, **kwargs): + # type: (PipelineType, **Any) -> Iterator[bytes] """Generator for streaming request body data.""" - return StreamDownloadGenerator(pipeline, self, decompress=decompress) + return StreamDownloadGenerator(pipeline, self, **kwargs) class RequestsTransport(HttpTransport): diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_trio.py b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_trio.py index 197a1529eff8..7762dfb2d518 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_trio.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_trio.py @@ -54,14 +54,15 @@ class TrioStreamDownloadGenerator(AsyncIterator): :param pipeline: The pipeline object :param response: The response object. - :param bool decompress: If True which is default, will attempt to decode the body based + :keyword bool decompress: If True which is default, will attempt to decode the body based on the ‘content-encoding’ header. """ - def __init__(self, pipeline: Pipeline, response: AsyncHttpResponse, decompress: bool = True) -> None: + def __init__(self, pipeline: Pipeline, response: AsyncHttpResponse, **kwargs) -> None: self.pipeline = pipeline self.request = response.request self.response = response self.block_size = response.block_size + decompress = kwargs.get("decompress", True) if decompress: self.iter_content_func = self.response.internal_response.iter_content(self.block_size) else: @@ -99,10 +100,10 @@ async def __anext__(self): class TrioRequestsTransportResponse(AsyncHttpResponse, RequestsTransportResponse): # type: ignore """Asynchronous streaming of data from the response. """ - def stream_download(self, pipeline, decompress=True) -> AsyncIteratorType[bytes]: # type: ignore + def stream_download(self, pipeline, **kwargs) -> AsyncIteratorType[bytes]: # type: ignore """Generator for streaming response data. """ - return TrioStreamDownloadGenerator(pipeline, self, decompress=decompress) + return TrioStreamDownloadGenerator(pipeline, self, **kwargs) class TrioRequestsTransport(RequestsAsyncTransportBase): # type: ignore From 85f0efee50c2e3aa7b241063f2853d20bbbd830d Mon Sep 17 00:00:00 2001 From: Xiang Yan Date: Wed, 5 May 2021 11:42:28 -0700 Subject: [PATCH 23/36] update --- sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py b/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py index f9ddc0829986..228f1e503ee6 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py @@ -300,7 +300,7 @@ def stream_download(self, pipeline, **kwargs) -> AsyncIteratorType[bytes]: :keyword bool decompress: If True which is default, will attempt to decode the body based on the ‘content-encoding’ header. """ - return AioHttpStreamDownloadGenerator(pipeline, self, decompress=decompress) + return AioHttpStreamDownloadGenerator(pipeline, self, **kwargs) def __getstate__(self): # Be sure body is loaded in memory, otherwise not pickable and let it throw From 510f93c7ad075bf32171d47cdd2bab7c863ebf6a Mon Sep 17 00:00:00 2001 From: Xiang Yan Date: Wed, 5 May 2021 11:52:04 -0700 Subject: [PATCH 24/36] update --- sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py b/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py index 228f1e503ee6..c9eaa524f81e 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py @@ -219,7 +219,7 @@ async def __anext__(self): raise _ResponseStopIteration() if not(self._decompress): return chunk - enc = response.internal_response.headers.get('Content-Encoding') + enc = self.response.internal_response.headers.get('Content-Encoding') if not(enc): return chunk enc = enc.lower() From fe523f72053bc6f76f326a8c31256d963adadd5f Mon Sep 17 00:00:00 2001 From: Xiang Yan Date: Wed, 5 May 2021 12:24:46 -0700 Subject: [PATCH 25/36] update --- .../azure-core/azure/core/pipeline/transport/_aiohttp.py | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py b/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py index c9eaa524f81e..71a4c9f9ff58 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py @@ -217,14 +217,14 @@ async def __anext__(self): chunk = await self.response.internal_response.content.read(self.block_size) if not chunk: raise _ResponseStopIteration() - if not(self._decompress): + if not self._decompress: return chunk enc = self.response.internal_response.headers.get('Content-Encoding') - if not(enc): + if not enc: return chunk enc = enc.lower() if enc in ("gzip", "deflate", "br"): - if encoding == "br": + if enc == "br": try: import brotli decompressor = brotli.Decompressor() @@ -236,7 +236,7 @@ async def __anext__(self): decompressor = brotli.Decompressor() else: import zlib - zlib_mode = 16 + zlib.MAX_WBITS if encoding == "gzip" else zlib.MAX_WBITS + zlib_mode = 16 + zlib.MAX_WBITS if enc == "gzip" else zlib.MAX_WBITS decompressor = zlib.decompressobj(wbits=zlib_mode) chunk = decompressor.decompress(chunk) return chunk From 8d4f4735c4300660e1b0749503aded3e1c6fea0f Mon Sep 17 00:00:00 2001 From: Xiang Yan Date: Wed, 5 May 2021 15:09:29 -0700 Subject: [PATCH 26/36] update --- sdk/core/azure-core/CLIENT_LIBRARY_DEVELOPER.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/sdk/core/azure-core/CLIENT_LIBRARY_DEVELOPER.md b/sdk/core/azure-core/CLIENT_LIBRARY_DEVELOPER.md index 7949e02f1627..7abd4fab832f 100644 --- a/sdk/core/azure-core/CLIENT_LIBRARY_DEVELOPER.md +++ b/sdk/core/azure-core/CLIENT_LIBRARY_DEVELOPER.md @@ -269,7 +269,7 @@ class HttpResponse(object): def text(self, encoding=None): """Return the whole body as a string.""" - def stream_download(self, pipeline, decompress=True): + def stream_download(self, pipeline, **kwargs): """Generator for streaming request body data. Should be implemented by sub-classes if streaming download is supported. From fa51573106d5fe7d884950886e89f39a02445574 Mon Sep 17 00:00:00 2001 From: Xiang Yan Date: Wed, 5 May 2021 18:56:33 -0700 Subject: [PATCH 27/36] update --- sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py | 4 +++- .../azure/core/pipeline/transport/_requests_asyncio.py | 4 +++- .../azure/core/pipeline/transport/_requests_basic.py | 4 +++- .../azure/core/pipeline/transport/_requests_trio.py | 4 +++- 4 files changed, 12 insertions(+), 4 deletions(-) diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py b/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py index 71a4c9f9ff58..abf6ba5a07d3 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py @@ -206,7 +206,9 @@ def __init__(self, pipeline: Pipeline, response: AsyncHttpResponse, **kwargs) -> self.request = response.request self.response = response self.block_size = response.block_size - self._decompress = kwargs.get("decompress", True) + self._decompress = kwargs.pop("decompress", True) + if len(kwargs) > 0: + raise ValueError("Unknown parameters!") self.content_length = int(response.internal_response.headers.get('Content-Length', 0)) def __len__(self): diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_asyncio.py b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_asyncio.py index 669be477b552..fe9eb06bf0ad 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_asyncio.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_asyncio.py @@ -146,7 +146,9 @@ def __init__(self, pipeline: Pipeline, response: AsyncHttpResponse, **kwargs) -> self.request = response.request self.response = response self.block_size = response.block_size - decompress = kwargs.get("decompress", True) + decompress = kwargs.pop("decompress", True) + if len(kwargs) > 0: + raise ValueError("Unknown parameters!") if decompress: self.iter_content_func = self.response.internal_response.iter_content(self.block_size) else: diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_basic.py b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_basic.py index 2d44ab0cd70c..ba55c58c1008 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_basic.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_basic.py @@ -128,7 +128,9 @@ def __init__(self, pipeline, response, **kwargs): self.request = response.request self.response = response self.block_size = response.block_size - decompress = kwargs.get("decompress", True) + decompress = kwargs.pop("decompress", True) + if len(kwargs) > 0: + raise ValueError("Unknown parameters!") if decompress: self.iter_content_func = self.response.internal_response.iter_content(self.block_size) else: diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_trio.py b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_trio.py index 7762dfb2d518..81f485e0f3c5 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_trio.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_trio.py @@ -62,7 +62,9 @@ def __init__(self, pipeline: Pipeline, response: AsyncHttpResponse, **kwargs) -> self.request = response.request self.response = response self.block_size = response.block_size - decompress = kwargs.get("decompress", True) + decompress = kwargs.pop("decompress", True) + if len(kwargs) > 0: + raise ValueError("Unknown parameters!") if decompress: self.iter_content_func = self.response.internal_response.iter_content(self.block_size) else: From b2e89e814cc3c9db7a59f3f039746765a891dd97 Mon Sep 17 00:00:00 2001 From: Xiang Yan Date: Wed, 5 May 2021 20:13:11 -0700 Subject: [PATCH 28/36] update --- .../azure/core/pipeline/transport/_aiohttp.py | 31 ++++++++++--------- 1 file changed, 16 insertions(+), 15 deletions(-) diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py b/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py index abf6ba5a07d3..047be4344664 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py @@ -210,6 +210,7 @@ def __init__(self, pipeline: Pipeline, response: AsyncHttpResponse, **kwargs) -> if len(kwargs) > 0: raise ValueError("Unknown parameters!") self.content_length = int(response.internal_response.headers.get('Content-Length', 0)) + self._decompressor = None def __len__(self): return self.content_length @@ -226,21 +227,21 @@ async def __anext__(self): return chunk enc = enc.lower() if enc in ("gzip", "deflate", "br"): - if enc == "br": - try: - import brotli - decompressor = brotli.Decompressor() - except ImportError: - raise DecodeError( - "Can not decode content-encoding: brotli (br). " - "Please install `brotlipy`" - ) - decompressor = brotli.Decompressor() - else: - import zlib - zlib_mode = 16 + zlib.MAX_WBITS if enc == "gzip" else zlib.MAX_WBITS - decompressor = zlib.decompressobj(wbits=zlib_mode) - chunk = decompressor.decompress(chunk) + if not self._decompressor: + if enc == "br": + try: + import brotli + self._decompressor = brotli.Decompressor() + except ImportError: + raise DecodeError( + "Can not decode content-encoding: brotli (br). " + "Please install `brotlipy`" + ) + else: + import zlib + zlib_mode = 16 + zlib.MAX_WBITS if enc == "gzip" else zlib.MAX_WBITS + self._decompressor = zlib.decompressobj(wbits=zlib_mode) + chunk = self._decompressor.decompress(chunk) return chunk except _ResponseStopIteration: self.response.internal_response.close() From 5c853351670360460655fe81175e2edca1cadb3e Mon Sep 17 00:00:00 2001 From: Xiang Yan Date: Wed, 5 May 2021 21:40:59 -0700 Subject: [PATCH 29/36] updates --- sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py | 2 +- .../azure/core/pipeline/transport/_requests_asyncio.py | 2 +- .../azure-core/azure/core/pipeline/transport/_requests_basic.py | 2 +- .../azure-core/azure/core/pipeline/transport/_requests_trio.py | 2 +- 4 files changed, 4 insertions(+), 4 deletions(-) diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py b/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py index 047be4344664..034d62e264a6 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py @@ -208,7 +208,7 @@ def __init__(self, pipeline: Pipeline, response: AsyncHttpResponse, **kwargs) -> self.block_size = response.block_size self._decompress = kwargs.pop("decompress", True) if len(kwargs) > 0: - raise ValueError("Unknown parameters!") + raise TypeError("Found incorrect parameter: {}".format(list(kwargs.keys())[0])) self.content_length = int(response.internal_response.headers.get('Content-Length', 0)) self._decompressor = None diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_asyncio.py b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_asyncio.py index fe9eb06bf0ad..66c37d1453d5 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_asyncio.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_asyncio.py @@ -148,7 +148,7 @@ def __init__(self, pipeline: Pipeline, response: AsyncHttpResponse, **kwargs) -> self.block_size = response.block_size decompress = kwargs.pop("decompress", True) if len(kwargs) > 0: - raise ValueError("Unknown parameters!") + raise TypeError("Found incorrect parameter: {}".format(list(kwargs.keys())[0])) if decompress: self.iter_content_func = self.response.internal_response.iter_content(self.block_size) else: diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_basic.py b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_basic.py index ba55c58c1008..1552c4731a18 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_basic.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_basic.py @@ -130,7 +130,7 @@ def __init__(self, pipeline, response, **kwargs): self.block_size = response.block_size decompress = kwargs.pop("decompress", True) if len(kwargs) > 0: - raise ValueError("Unknown parameters!") + raise TypeError("Found incorrect parameter: {}".format(list(kwargs.keys())[0])) if decompress: self.iter_content_func = self.response.internal_response.iter_content(self.block_size) else: diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_trio.py b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_trio.py index 81f485e0f3c5..c48131c1b410 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_trio.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_trio.py @@ -64,7 +64,7 @@ def __init__(self, pipeline: Pipeline, response: AsyncHttpResponse, **kwargs) -> self.block_size = response.block_size decompress = kwargs.pop("decompress", True) if len(kwargs) > 0: - raise ValueError("Unknown parameters!") + raise TypeError("Found incorrect parameter: {}".format(list(kwargs.keys())[0])) if decompress: self.iter_content_func = self.response.internal_response.iter_content(self.block_size) else: From c5decbc1e49236e39e360276ba08d9a9235c1abd Mon Sep 17 00:00:00 2001 From: Xiang Yan Date: Wed, 5 May 2021 21:57:58 -0700 Subject: [PATCH 30/36] update --- sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py | 2 +- .../azure/core/pipeline/transport/_requests_asyncio.py | 2 +- .../azure-core/azure/core/pipeline/transport/_requests_basic.py | 2 +- .../azure-core/azure/core/pipeline/transport/_requests_trio.py | 2 +- 4 files changed, 4 insertions(+), 4 deletions(-) diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py b/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py index 034d62e264a6..779faa5623d6 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py @@ -208,7 +208,7 @@ def __init__(self, pipeline: Pipeline, response: AsyncHttpResponse, **kwargs) -> self.block_size = response.block_size self._decompress = kwargs.pop("decompress", True) if len(kwargs) > 0: - raise TypeError("Found incorrect parameter: {}".format(list(kwargs.keys())[0])) + raise TypeError("Got an unexpected keyword argument: {}".format(list(kwargs.keys())[0])) self.content_length = int(response.internal_response.headers.get('Content-Length', 0)) self._decompressor = None diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_asyncio.py b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_asyncio.py index 66c37d1453d5..4f070834fa85 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_asyncio.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_asyncio.py @@ -148,7 +148,7 @@ def __init__(self, pipeline: Pipeline, response: AsyncHttpResponse, **kwargs) -> self.block_size = response.block_size decompress = kwargs.pop("decompress", True) if len(kwargs) > 0: - raise TypeError("Found incorrect parameter: {}".format(list(kwargs.keys())[0])) + raise TypeError("Got an unexpected keyword argument: {}".format(list(kwargs.keys())[0])) if decompress: self.iter_content_func = self.response.internal_response.iter_content(self.block_size) else: diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_basic.py b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_basic.py index 1552c4731a18..4cba767842ff 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_basic.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_basic.py @@ -130,7 +130,7 @@ def __init__(self, pipeline, response, **kwargs): self.block_size = response.block_size decompress = kwargs.pop("decompress", True) if len(kwargs) > 0: - raise TypeError("Found incorrect parameter: {}".format(list(kwargs.keys())[0])) + raise TypeError("Got an unexpected keyword argument: {}".format(list(kwargs.keys())[0])) if decompress: self.iter_content_func = self.response.internal_response.iter_content(self.block_size) else: diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_trio.py b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_trio.py index c48131c1b410..58fa2722e2d0 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_requests_trio.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_requests_trio.py @@ -64,7 +64,7 @@ def __init__(self, pipeline: Pipeline, response: AsyncHttpResponse, **kwargs) -> self.block_size = response.block_size decompress = kwargs.pop("decompress", True) if len(kwargs) > 0: - raise TypeError("Found incorrect parameter: {}".format(list(kwargs.keys())[0])) + raise TypeError("Got an unexpected keyword argument: {}".format(list(kwargs.keys())[0])) if decompress: self.iter_content_func = self.response.internal_response.iter_content(self.block_size) else: From 66d0c5509b88d5894719056345406bc61844e754 Mon Sep 17 00:00:00 2001 From: Xiang Yan Date: Wed, 5 May 2021 22:26:35 -0700 Subject: [PATCH 31/36] update --- .../azure/core/pipeline/transport/_aiohttp.py | 23 ++++++++++++++++--- 1 file changed, 20 insertions(+), 3 deletions(-) diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py b/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py index 779faa5623d6..742f6d4fea79 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py @@ -193,6 +193,24 @@ async def send(self, request: HttpRequest, **config: Any) -> Optional[AsyncHttpR return response +class _BrotliDecoder: + # Supports both 'brotlipy' and 'Brotli' packages + # since they share an import name. The top branches + # are for 'brotlipy' and bottom branches for 'Brotli' + def __init__(self) -> None: + import brotli + self._obj = brotli.Decompressor() + + def decompress(self, data: bytes) -> bytes: + if hasattr(self._obj, "decompress"): + return cast(bytes, self._obj.decompress(data)) + return cast(bytes, self._obj.process(data)) + + def flush(self) -> bytes: + if hasattr(self._obj, "flush"): + return cast(bytes, self._obj.flush()) + return b"" + class AioHttpStreamDownloadGenerator(AsyncIterator): """Streams the response body data. @@ -230,12 +248,11 @@ async def __anext__(self): if not self._decompressor: if enc == "br": try: - import brotli - self._decompressor = brotli.Decompressor() + self._decompressor = _BrotliDecoder() except ImportError: raise DecodeError( "Can not decode content-encoding: brotli (br). " - "Please install `brotlipy`" + "Please install `Brotli`" ) else: import zlib From 3853b8ae4b9e936c4f93118ac7bfe0a40a34d939 Mon Sep 17 00:00:00 2001 From: Xiang Yan Date: Wed, 5 May 2021 22:35:04 -0700 Subject: [PATCH 32/36] update --- .../azure/core/pipeline/transport/_aiohttp.py | 17 ++++------------- 1 file changed, 4 insertions(+), 13 deletions(-) diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py b/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py index 742f6d4fea79..dfc7a753f6a4 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py @@ -244,20 +244,11 @@ async def __anext__(self): if not enc: return chunk enc = enc.lower() - if enc in ("gzip", "deflate", "br"): + if enc in ("gzip", "deflate"): if not self._decompressor: - if enc == "br": - try: - self._decompressor = _BrotliDecoder() - except ImportError: - raise DecodeError( - "Can not decode content-encoding: brotli (br). " - "Please install `Brotli`" - ) - else: - import zlib - zlib_mode = 16 + zlib.MAX_WBITS if enc == "gzip" else zlib.MAX_WBITS - self._decompressor = zlib.decompressobj(wbits=zlib_mode) + import zlib + zlib_mode = 16 + zlib.MAX_WBITS if enc == "gzip" else zlib.MAX_WBITS + self._decompressor = zlib.decompressobj(wbits=zlib_mode) chunk = self._decompressor.decompress(chunk) return chunk except _ResponseStopIteration: From 338b46b89b09114bd6c559554c325c393cc61751 Mon Sep 17 00:00:00 2001 From: Xiang Yan Date: Thu, 6 May 2021 08:53:45 -0700 Subject: [PATCH 33/36] update --- .../azure/core/pipeline/transport/_aiohttp.py | 18 ------------------ 1 file changed, 18 deletions(-) diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py b/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py index dfc7a753f6a4..098da9e98209 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py @@ -193,24 +193,6 @@ async def send(self, request: HttpRequest, **config: Any) -> Optional[AsyncHttpR return response -class _BrotliDecoder: - # Supports both 'brotlipy' and 'Brotli' packages - # since they share an import name. The top branches - # are for 'brotlipy' and bottom branches for 'Brotli' - def __init__(self) -> None: - import brotli - self._obj = brotli.Decompressor() - - def decompress(self, data: bytes) -> bytes: - if hasattr(self._obj, "decompress"): - return cast(bytes, self._obj.decompress(data)) - return cast(bytes, self._obj.process(data)) - - def flush(self) -> bytes: - if hasattr(self._obj, "flush"): - return cast(bytes, self._obj.flush()) - return b"" - class AioHttpStreamDownloadGenerator(AsyncIterator): """Streams the response body data. From 9f40bbea30062df4771ff89607d1cdfd45952b17 Mon Sep 17 00:00:00 2001 From: Xiang Yan Date: Thu, 6 May 2021 09:11:16 -0700 Subject: [PATCH 34/36] update --- sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py | 1 - 1 file changed, 1 deletion(-) diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py b/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py index 098da9e98209..b25944b9b0d3 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py @@ -192,7 +192,6 @@ async def send(self, request: HttpRequest, **config: Any) -> Optional[AsyncHttpR raise ServiceResponseError(err, error=err) from err return response - class AioHttpStreamDownloadGenerator(AsyncIterator): """Streams the response body data. From 793626229f2f3c2cbd2e9fe93f4a03e4f627e1d2 Mon Sep 17 00:00:00 2001 From: Xiang Yan Date: Thu, 6 May 2021 09:40:48 -0700 Subject: [PATCH 35/36] update --- sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py b/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py index b25944b9b0d3..2a2641b90dbc 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py @@ -33,7 +33,7 @@ from requests.exceptions import StreamConsumedError from azure.core.configuration import ConnectionConfiguration -from azure.core.exceptions import ServiceRequestError, ServiceResponseError, DecodeError +from azure.core.exceptions import ServiceRequestError, ServiceResponseError from azure.core.pipeline import Pipeline from ._base import HttpRequest From 2616dd7ca3e17c76c97febc328c8656463b40882 Mon Sep 17 00:00:00 2001 From: Xiang Yan Date: Thu, 6 May 2021 10:11:10 -0700 Subject: [PATCH 36/36] update --- sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py | 1 - 1 file changed, 1 deletion(-) diff --git a/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py b/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py index 2a2641b90dbc..da70db52b5be 100644 --- a/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py +++ b/sdk/core/azure-core/azure/core/pipeline/transport/_aiohttp.py @@ -46,7 +46,6 @@ CONTENT_CHUNK_SIZE = 10 * 1024 _LOGGER = logging.getLogger(__name__) - class AioHttpTransport(AsyncHttpTransport): """AioHttp HTTP sender implementation.