Team Ai
Datasetpublic

codekingpro/portable-devtools

sourceHugging Faceupdated 5mo agoView on Hugging Face
1likes15kdownloads
payload_streamer.py79 linesDownload Raw Back to aiohttp
1"""
2Payload implementation for coroutines as data provider.
3
4As a simple case, you can upload data from file::
5
6   @aiohttp.streamer
7   async def file_sender(writer, file_name=None):
8      with open(file_name, 'rb') as f:
9          chunk = f.read(2**16)
10          while chunk:
11              await writer.write(chunk)
12
13              chunk = f.read(2**16)
14
15Then you can use `file_sender` like this:
16
17    async with session.post('http://httpbin.org/post',
18                            data=file_sender(file_name='huge_file')) as resp:
19        print(await resp.text())
20
21..note:: Coroutine must accept `writer` as first argument
22
23"""
24
25import types
26import warnings
27from typing import Any, Awaitable, Callable, Dict, Tuple
28
29from .abc import AbstractStreamWriter
30from .payload import Payload, payload_type
31
32__all__ = ("streamer",)
33
34
35class _stream_wrapper:
36    def __init__(
37        self,
38        coro: Callable[..., Awaitable[None]],
39        args: Tuple[Any, ...],
40        kwargs: Dict[str, Any],
41    ) -> None:
42        self.coro = types.coroutine(coro)
43        self.args = args
44        self.kwargs = kwargs
45
46    async def __call__(self, writer: AbstractStreamWriter) -> None:
47        await self.coro(writer, *self.args, **self.kwargs)
48
49
50class streamer:
51    def __init__(self, coro: Callable[..., Awaitable[None]]) -> None:
52        warnings.warn(
53            "@streamer is deprecated, use async generators instead",
54            DeprecationWarning,
55            stacklevel=2,
56        )
57        self.coro = coro
58
59    def __call__(self, *args: Any, **kwargs: Any) -> _stream_wrapper:
60        return _stream_wrapper(self.coro, args, kwargs)
61
62
63@payload_type(_stream_wrapper)
64class StreamWrapperPayload(Payload):
65    async def write(self, writer: AbstractStreamWriter) -> None:
66        await self._value(writer)
67
68    def decode(self, encoding: str = "utf-8", errors: str = "strict") -> str:
69        raise TypeError("Unable to decode.")
70
71
72@payload_type(streamer)
73class StreamPayload(StreamWrapperPayload):
74    def __init__(self, value: Any, *args: Any, **kwargs: Any) -> None:
75        super().__init__(value(), *args, **kwargs)
76
77    async def write(self, writer: AbstractStreamWriter) -> None:
78        await self._value(writer)
79 
codekingpro/portable-devtools · Team Ai