Team Ai
Datasetpublic

codekingpro/portable-devtools

sourceHugging Faceupdated 5mo agoView on Hugging Face
1likes14kdownloads
async_abc.py165 linesDownload Raw Back to pipeline
1# --------------------------------------------------------------------------2#3# Copyright (c) Microsoft Corporation. All rights reserved.4#5# The MIT License (MIT)6#7# Permission is hereby granted, free of charge, to any person obtaining a copy8# of this software and associated documentation files (the ""Software""), to9# deal in the Software without restriction, including without limitation the10# rights to use, copy, modify, merge, publish, distribute, sublicense, and/or11# sell copies of the Software, and to permit persons to whom the Software is12# furnished to do so, subject to the following conditions:13#14# The above copyright notice and this permission notice shall be included in15# all copies or substantial portions of the Software.16#17# THE SOFTWARE IS PROVIDED *AS IS*, WITHOUT WARRANTY OF ANY KIND, EXPRESS OR18# IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY,19# FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE20# AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER21# LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING22# FROM, OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS23# IN THE SOFTWARE.24#25# --------------------------------------------------------------------------26import abc27 28from typing import Any, List, Union, Callable, AsyncIterator, Optional, Generic, TypeVar29 30from . import Request, Response, Pipeline, SansIOHTTPPolicy31 32 33AsyncHTTPResponseType = TypeVar("AsyncHTTPResponseType")34HTTPRequestType = TypeVar("HTTPRequestType")35 36try:37    from contextlib import AbstractAsyncContextManager  # type: ignore38except ImportError: # Python <= 3.739    class AbstractAsyncContextManager(object):  # type: ignore40        async def __aenter__(self):41            """Return `self` upon entering the runtime context."""42            return self43 44        @abc.abstractmethod45        async def __aexit__(self, exc_type, exc_value, traceback):46            """Raise any exception triggered within the runtime context."""47            return None48 49 50 51 52class AsyncHTTPPolicy(abc.ABC, Generic[HTTPRequestType, AsyncHTTPResponseType]):53    """An http policy ABC.54    """55    def __init__(self) -> None:56        # next will be set once in the pipeline57        self.next = None  # type: Optional[Union[AsyncHTTPPolicy[HTTPRequestType, AsyncHTTPResponseType], AsyncHTTPSender[HTTPRequestType, AsyncHTTPResponseType]]]58 59    @abc.abstractmethod60    async def send(self, request: Request, **kwargs: Any) -> Response[HTTPRequestType, AsyncHTTPResponseType]:61        """Mutate the request.62 63        Context content is dependent of the HTTPSender.64        """65        pass66 67 68class _SansIOAsyncHTTPPolicyRunner(AsyncHTTPPolicy[HTTPRequestType, AsyncHTTPResponseType]):69    """Async implementation of the SansIO policy.70    """71 72    def __init__(self, policy: SansIOHTTPPolicy) -> None:73        super(_SansIOAsyncHTTPPolicyRunner, self).__init__()74        self._policy = policy75 76    async def send(self, request: Request, **kwargs: Any) -> Response[HTTPRequestType, AsyncHTTPResponseType]:77        self._policy.on_request(request, **kwargs)78        try:79            response = await self.next.send(request, **kwargs)  # type: ignore80        except Exception:81            if not self._policy.on_exception(request, **kwargs):82                raise83        else:84            self._policy.on_response(request, response, **kwargs)85        return response86 87 88class AsyncHTTPSender(AbstractAsyncContextManager, abc.ABC, Generic[HTTPRequestType, AsyncHTTPResponseType]):89    """An http sender ABC.90    """91 92    @abc.abstractmethod93    async def send(self, request: Request[HTTPRequestType], **config: Any) -> Response[HTTPRequestType, AsyncHTTPResponseType]:94        """Send the request using this HTTP sender.95        """96        pass97 98    def build_context(self) -> Any:99        """Allow the sender to build a context that will be passed100        across the pipeline with the request.101 102        Return type has no constraints. Implementation is not103        required and None by default.104        """105        return None106 107    def __enter__(self):108        raise TypeError("Use async with instead")109 110    def __exit__(self, exc_type, exc_val, exc_tb):111        # __exit__ should exist in pair with __enter__ but never executed112        pass  # pragma: no cover113 114 115class AsyncPipeline(AbstractAsyncContextManager, Generic[HTTPRequestType, AsyncHTTPResponseType]):116    """A pipeline implementation.117 118    This is implemented as a context manager, that will activate the context119    of the HTTP sender.120    """121 122    def __init__(self, policies: List[Union[AsyncHTTPPolicy, SansIOHTTPPolicy]] = None, sender: Optional[AsyncHTTPSender[HTTPRequestType, AsyncHTTPResponseType]] = None) -> None:123        self._impl_policies = []  # type: List[AsyncHTTPPolicy[HTTPRequestType, AsyncHTTPResponseType]]124        if sender:125            self._sender = sender126        else:127            # Import default only if nothing is provided128            from .aiohttp import AioHTTPSender129            self._sender = AioHTTPSender()130 131        for policy in (policies or []):132            if isinstance(policy, SansIOHTTPPolicy):133                self._impl_policies.append(_SansIOAsyncHTTPPolicyRunner(policy))134            else:135                self._impl_policies.append(policy)136        for index in range(len(self._impl_policies)-1):137            self._impl_policies[index].next = self._impl_policies[index+1]138        if self._impl_policies:139            self._impl_policies[-1].next = self._sender140 141    def __enter__(self):142        raise TypeError("Use 'async with' instead")143 144    def __exit__(self, exc_type, exc_val, exc_tb):145        # __exit__ should exist in pair with __enter__ but never executed146        pass  # pragma: no cover147 148    async def __aenter__(self) -> 'AsyncPipeline':149        await self._sender.__aenter__()150        return self151 152    async def __aexit__(self, *exc_details):  # pylint: disable=arguments-differ153        await self._sender.__aexit__(*exc_details)154 155    async def run(self, request: Request, **kwargs: Any) -> Response[HTTPRequestType, AsyncHTTPResponseType]:156        context = self._sender.build_context()157        pipeline_request = Request(request, context)158        first_node = self._impl_policies[0] if self._impl_policies else self._sender159        return await first_node.send(pipeline_request, **kwargs)  # type: ignore160 161__all__ = [162    'AsyncHTTPPolicy',163    'AsyncHTTPSender',164    'AsyncPipeline',165]
codekingpro/portable-devtools · Team Ai