/
opt
/
cloudlinux
/
venv
/
lib
/
python3.11
/
site-packages
/
aiohttp
/
/opt/cloudlinux/venv/lib/python3.11/site-packages/aiohttp
mkdir
upload
Name
Size
Mode
Actions
.hash/
-
0755
rm
__pycache__/
-
0755
rm
abc.py
5500
0644
edit
dl
rm
base_protocol.py
2741
0644
edit
dl
rm
client.py
47276
0644
edit
dl
rm
client_exceptions.py
9411
0644
edit
dl
rm
client_proto.py
8651
0644
edit
dl
rm
client_reqrep.py
39680
0644
edit
dl
rm
client_ws.py
11010
0644
edit
dl
rm
compression_utils.py
5015
0644
edit
dl
rm
connector.py
52798
0644
edit
dl
rm
cookiejar.py
14015
0644
edit
dl
rm
formdata.py
6106
0644
edit
dl
rm
hdrs.py
4613
0644
edit
dl
rm
helpers.py
30255
0644
edit
dl
rm
http.py
1842
0644
edit
dl
rm
http_exceptions.py
2716
0644
edit
dl
rm
http_parser.py
35496
0644
edit
dl
rm
http_websocket.py
26716
0644
edit
dl
rm
http_writer.py
5933
0644
edit
dl
rm
locks.py
1136
0644
edit
dl
rm
log.py
325
0644
edit
dl
rm
multipart.py
32472
0644
edit
dl
rm
payload.py
13542
0644
edit
dl
rm
payload_streamer.py
2087
0644
edit
dl
rm
py.typed
7
0644
edit
dl
rm
pytest_plugin.py
11605
0644
edit
dl
rm
resolver.py
5070
0644
edit
dl
rm
streams.py
20836
0644
edit
dl
rm
tcp_helpers.py
961
0644
edit
dl
rm
test_utils.py
20185
0644
edit
dl
rm
tracing.py
15132
0644
edit
dl
rm
typedefs.py
1471
0644
edit
dl
rm
web.py
19263
0644
edit
dl
rm
web_app.py
18311
0644
edit
dl
rm
web_exceptions.py
10360
0644
edit
dl
rm
web_fileresponse.py
11416
0644
edit
dl
rm
web_log.py
7801
0644
edit
dl
rm
web_middlewares.py
4032
0644
edit
dl
rm
web_protocol.py
23044
0644
edit
dl
rm
web_request.py
28756
0644
edit
dl
rm
web_response.py
27729
0644
edit
dl
rm
web_routedef.py
6132
0644
edit
dl
rm
web_runner.py
11736
0644
edit
dl
rm
web_server.py
2587
0644
edit
dl
rm
web_urldispatcher.py
40057
0644
edit
dl
rm
web_ws.py
18647
0644
edit
dl
rm
worker.py
7965
0644
edit
dl
rm
_cparser.pxd
4318
0644
edit
dl
rm
_find_header.pxd
68
0644
edit
dl
rm
_headers.pxi
2007
0644
edit
dl
rm
_helpers.cpython-311-x86_64-linux-gnu.so
310476
0755
edit
dl
rm
_helpers.pyi
202
0644
edit
dl
rm
_helpers.pyx
1049
0644
edit
dl
rm
_http_parser.cpython-311-x86_64-linux-gnu.so
1889106
0755
edit
dl
rm
_http_parser.pyx
28058
0644
edit
dl
rm
_http_writer.cpython-311-x86_64-linux-gnu.so
320845
0755
edit
dl
rm
_http_writer.pyx
4575
0644
edit
dl
rm
_websocket.cpython-311-x86_64-linux-gnu.so
187629
0755
edit
dl
rm
_websocket.pyx
1561
0644
edit
dl
rm
__init__.py
7762
0644
edit
dl
rm
Edit:
/opt/cloudlinux/venv/lib/python3.11/site-packages/aiohttp/client_proto.py
(8651B)
import asyncio from contextlib import suppress from typing import Any, Optional, Tuple from .base_protocol import BaseProtocol from .client_exceptions import ( ClientOSError, ClientPayloadError, ServerDisconnectedError, ServerTimeoutError, ) from .helpers import BaseTimerContext, status_code_must_be_empty_body from .http import HttpResponseParser, RawResponseMessage from .streams import EMPTY_PAYLOAD, DataQueue, StreamReader class ResponseHandler(BaseProtocol, DataQueue[Tuple[RawResponseMessage, StreamReader]]): """Helper class to adapt between Protocol and StreamReader.""" def __init__(self, loop: asyncio.AbstractEventLoop) -> None: BaseProtocol.__init__(self, loop=loop) DataQueue.__init__(self, loop) self._should_close = False self._payload: Optional[StreamReader] = None self._skip_payload = False self._payload_parser = None self._timer = None self._tail = b"" self._upgraded = False self._parser: Optional[HttpResponseParser] = None self._read_timeout: Optional[float] = None self._read_timeout_handle: Optional[asyncio.TimerHandle] = None self._timeout_ceil_threshold: Optional[float] = 5 @property def upgraded(self) -> bool: return self._upgraded @property def should_close(self) -> bool: if self._payload is not None and not self._payload.is_eof() or self._upgraded: return True return ( self._should_close or self._upgraded or self.exception() is not None or self._payload_parser is not None or len(self) > 0 or bool(self._tail) ) def force_close(self) -> None: self._should_close = True def close(self) -> None: transport = self.transport if transport is not None: transport.close() self.transport = None self._payload = None self._drop_timeout() def is_connected(self) -> bool: return self.transport is not None and not self.transport.is_closing() def connection_lost(self, exc: Optional[BaseException]) -> None: self._drop_timeout() if self._payload_parser is not None: with suppress(Exception): self._payload_parser.feed_eof() uncompleted = None if self._parser is not None: try: uncompleted = self._parser.feed_eof() except Exception as e: if self._payload is not None: exc = ClientPayloadError("Response payload is not completed") exc.__cause__ = e self._payload.set_exception(exc) if not self.is_eof(): if isinstance(exc, OSError): exc = ClientOSError(*exc.args) if exc is None: exc = ServerDisconnectedError(uncompleted) # assigns self._should_close to True as side effect, # we do it anyway below self.set_exception(exc) self._should_close = True self._parser = None self._payload = None self._payload_parser = None self._reading_paused = False super().connection_lost(exc) def eof_received(self) -> None: # should call parser.feed_eof() most likely self._drop_timeout() def pause_reading(self) -> None: super().pause_reading() self._drop_timeout() def resume_reading(self) -> None: super().resume_reading() self._reschedule_timeout() def set_exception(self, exc: BaseException) -> None: self._should_close = True self._drop_timeout() super().set_exception(exc) def set_parser(self, parser: Any, payload: Any) -> None: # TODO: actual types are: # parser: WebSocketReader # payload: FlowControlDataQueue # but they are not generi enough # Need an ABC for both types self._payload = payload self._payload_parser = parser self._drop_timeout() if self._tail: data, self._tail = self._tail, b"" self.data_received(data) def set_response_params( self, *, timer: Optional[BaseTimerContext] = None, skip_payload: bool = False, read_until_eof: bool = False, auto_decompress: bool = True, read_timeout: Optional[float] = None, read_bufsize: int = 2**16, timeout_ceil_threshold: float = 5, max_line_size: int = 8190, max_field_size: int = 8190, ) -> None: self._skip_payload = skip_payload self._read_timeout = read_timeout self._timeout_ceil_threshold = timeout_ceil_threshold self._parser = HttpResponseParser( self, self._loop, read_bufsize, timer=timer, payload_exception=ClientPayloadError, response_with_body=not skip_payload, read_until_eof=read_until_eof, auto_decompress=auto_decompress, max_line_size=max_line_size, max_field_size=max_field_size, ) if self._tail: data, self._tail = self._tail, b"" self.data_received(data) def _drop_timeout(self) -> None: if self._read_timeout_handle is not None: self._read_timeout_handle.cancel() self._read_timeout_handle = None def _reschedule_timeout(self) -> None: timeout = self._read_timeout if self._read_timeout_handle is not None: self._read_timeout_handle.cancel() if timeout: self._read_timeout_handle = self._loop.call_later( timeout, self._on_read_timeout ) else: self._read_timeout_handle = None def start_timeout(self) -> None: self._reschedule_timeout() def _on_read_timeout(self) -> None: exc = ServerTimeoutError("Timeout on reading data from socket") self.set_exception(exc) if self._payload is not None: self._payload.set_exception(exc) def data_received(self, data: bytes) -> None: self._reschedule_timeout() if not data: return # custom payload parser if self._payload_parser is not None: eof, tail = self._payload_parser.feed_data(data) if eof: self._payload = None self._payload_parser = None if tail: self.data_received(tail) return else: if self._upgraded or self._parser is None: # i.e. websocket connection, websocket parser is not set yet self._tail += data else: # parse http messages try: messages, upgraded, tail = self._parser.feed_data(data) except BaseException as exc: if self.transport is not None: # connection.release() could be called BEFORE # data_received(), the transport is already # closed in this case self.transport.close() # should_close is True after the call self.set_exception(exc) return self._upgraded = upgraded payload: Optional[StreamReader] = None for message, payload in messages: if message.should_close: self._should_close = True self._payload = payload if self._skip_payload or status_code_must_be_empty_body( message.code ): self.feed_data((message, EMPTY_PAYLOAD), 0) else: self.feed_data((message, payload), 0) if payload is not None: # new message(s) was processed # register timeout handler unsubscribing # either on end-of-stream or immediately for # EMPTY_PAYLOAD if payload is not EMPTY_PAYLOAD: payload.on_eof(self._drop_timeout) else: self._drop_timeout() if tail: if upgraded: self.data_received(tail) else: self._tail = tail
Save
cmd:
run