import asyncio from contextlib import suppress from typing import Any, Optional, Tuple, Union
from .base_protocol import BaseProtocol from .client_exceptions import (
ClientConnectionError,
ClientOSError,
ClientPayloadError,
ServerDisconnectedError,
SocketTimeoutError,
) from .helpers import (
_EXC_SENTINEL,
EMPTY_BODY_STATUS_CODES,
BaseTimerContext,
set_exception,
set_result,
) from .http import HttpResponseParser, RawResponseMessage from .http_exceptions import HttpProcessingError from .streams import EMPTY_PAYLOAD, DataQueue, StreamReader
class ResponseHandler(BaseProtocol, DataQueue[Tuple[RawResponseMessage, StreamReader]]): """Helper class to adapt between Protocol and StreamReader."""
@property def closed(self) -> Union[None, asyncio.Future[None]]: """Future that is set when the connection is closed.
This property returns a Future that will be completed when the connection is closed. The Future is created lazily on first access to avoid creating
futures that will never be awaited.
Returns:
- A Future[None] if the connection is still open or was closed after
this property was accessed
- Noneif connection_lost() was already called before this property
was ever accessed (indicating no one is waiting for the closure) """ if self._closed isNoneandnot self._connection_lost_called:
self._closed = self._loop.create_future() return self._closed
@property def should_close(self) -> bool: return bool(
self._should_close or (self._payload isnotNoneandnot self._payload.is_eof()) or self._upgraded or self._exception isnotNone or self._payload_parser isnotNone or self._buffer or self._tail
)
if self._closed isnotNone: # If someone is waiting for the closed future, # we should set it to None or an exception. If # self._closed is None, it means that # connection_lost() was called already # or nobody is waiting for it. if connection_closed_cleanly:
set_result(self._closed, None) else: assert original_connection_error isnotNone
set_exception(
self._closed,
ClientConnectionError(
f"Connection lost: {original_connection_error !s}",
),
original_connection_error,
)
if self._payload_parser isnotNone: with suppress(Exception): # FIXME: log this somehow?
self._payload_parser.feed_eof()
uncompleted = None if self._parser isnotNone: try:
uncompleted = self._parser.feed_eof() except Exception as underlying_exc: if self._payload isnotNone:
client_payload_exc_msg = (
f"Response payload is not completed: {underlying_exc !r}"
) ifnot connection_closed_cleanly:
client_payload_exc_msg = (
f"{client_payload_exc_msg !s}. "
f"{original_connection_error !r}"
)
set_exception(
self._payload,
ClientPayloadError(client_payload_exc_msg),
underlying_exc,
)
ifnot self.is_eof(): if isinstance(original_connection_error, OSError):
reraised_exc = ClientOSError(*original_connection_error.args) if connection_closed_cleanly:
reraised_exc = ServerDisconnectedError(uncompleted) # assigns self._should_close to True as side effect, # we do it anyway below
underlying_non_eof_exc = (
_EXC_SENTINEL if connection_closed_cleanly else original_connection_error
) assert underlying_non_eof_exc isnotNone assert reraised_exc isnotNone
self.set_exception(reraised_exc, underlying_non_eof_exc)
def set_parser(self, parser: Any, payload: Any) -> None: # TODO: actual types are: # parser: WebSocketReader # payload: WebSocketDataQueue # 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 _on_read_timeout(self) -> None:
exc = SocketTimeoutError("Timeout on reading data from socket")
self.set_exception(exc) if self._payload isnotNone:
set_exception(self._payload, exc)
# custom payload parser - currently always WebSocketReader if self._payload_parser isnotNone:
eof, tail = self._payload_parser.feed_data(data) if eof:
self._payload = None
self._payload_parser = None
if tail:
self.data_received(tail) return
if self._upgraded or self._parser isNone: # i.e. websocket connection, websocket parser is not set yet
self._tail += data return
# parse http messages try:
messages, upgraded, tail = self._parser.feed_data(data) except BaseException as underlying_exc: if self.transport isnotNone: # 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 if isinstance(underlying_exc, HttpProcessingError):
exc = HttpProcessingError(
code=underlying_exc.code,
message=underlying_exc.message,
headers=underlying_exc.headers,
) else:
exc = HttpProcessingError()
self.set_exception(exc, underlying_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 message.code in EMPTY_BODY_STATUS_CODES:
self.feed_data((message, EMPTY_PAYLOAD), 0) else:
self.feed_data((message, payload), 0)
if payload isnotNone: # new message(s) was processed # register timeout handler unsubscribing # either on end-of-stream or immediately for # EMPTY_PAYLOAD if payload isnot EMPTY_PAYLOAD:
payload.on_eof(self._drop_timeout) else:
self._drop_timeout()
if upgraded and tail:
self.data_received(tail)
Messung V0.5 in Prozent
¤ Dauer der Verarbeitung: 0.22 Sekunden
(vorverarbeitet am 2026-08-25)
¤
Die Informationen auf dieser Webseite wurden
nach bestem Wissen sorgfältig zusammengestellt. Es wird jedoch weder Vollständigkeit, noch Richtigkeit,
noch Qualität der bereit gestellten Informationen zugesichert.
Bemerkung:
Die farbliche Syntaxdarstellung und die Messung sind noch experimentell.