Team Ai
Datasetpublic

codekingpro/portable-devtools

sourceHugging Faceupdated 5mo agoView on Hugging Face
1likes14kdownloads
connection.py2111 linesDownload Raw Back to h2
1"""2h2/connection3~~~~~~~~~~~~~4 5An implementation of a HTTP/2 connection.6"""7from __future__ import annotations8 9import base6410from enum import Enum, IntEnum11from typing import TYPE_CHECKING, Any, Callable12 13from hpack.exceptions import HPACKError, OversizedHeaderListError14from hpack.hpack import Decoder, Encoder15from hyperframe.exceptions import InvalidPaddingError16from hyperframe.frame import (17    AltSvcFrame,18    ContinuationFrame,19    DataFrame,20    ExtensionFrame,21    Frame,22    GoAwayFrame,23    HeadersFrame,24    PingFrame,25    PriorityFrame,26    PushPromiseFrame,27    RstStreamFrame,28    SettingsFrame,29    WindowUpdateFrame,30)31 32from .config import H2Configuration33from .errors import ErrorCodes, _error_code_from_int34from .events import (35    AlternativeServiceAvailable,36    ConnectionTerminated,37    Event,38    InformationalResponseReceived,39    PingAckReceived,40    PingReceived,41    PriorityUpdated,42    RemoteSettingsChanged,43    RequestReceived,44    ResponseReceived,45    SettingsAcknowledged,46    TrailersReceived,47    UnknownFrameReceived,48    WindowUpdated,49)50from .exceptions import (51    DenialOfServiceError,52    FlowControlError,53    FrameTooLargeError,54    NoAvailableStreamIDError,55    NoSuchStreamError,56    ProtocolError,57    RFC1122Error,58    StreamClosedError,59    StreamIDTooLowError,60    TooManyStreamsError,61)62from .frame_buffer import FrameBuffer63from .settings import ChangedSetting, SettingCodes, Settings64from .stream import H2Stream, StreamClosedBy65from .utilities import SizeLimitDict, guard_increment_window66from .windows import WindowManager67 68if TYPE_CHECKING:  # pragma: no cover69    from collections.abc import Iterable70 71    from hpack.struct import Header, HeaderWeaklyTyped72 73 74class ConnectionState(Enum):75    IDLE = 076    CLIENT_OPEN = 177    SERVER_OPEN = 278    CLOSED = 379 80 81class ConnectionInputs(Enum):82    SEND_HEADERS = 083    SEND_PUSH_PROMISE = 184    SEND_DATA = 285    SEND_GOAWAY = 386    SEND_WINDOW_UPDATE = 487    SEND_PING = 588    SEND_SETTINGS = 689    SEND_RST_STREAM = 790    SEND_PRIORITY = 891    RECV_HEADERS = 992    RECV_PUSH_PROMISE = 1093    RECV_DATA = 1194    RECV_GOAWAY = 1295    RECV_WINDOW_UPDATE = 1396    RECV_PING = 1497    RECV_SETTINGS = 1598    RECV_RST_STREAM = 1699    RECV_PRIORITY = 17100    SEND_ALTERNATIVE_SERVICE = 18  # Added in 2.3.0101    RECV_ALTERNATIVE_SERVICE = 19  # Added in 2.3.0102 103 104class AllowedStreamIDs(IntEnum):105    EVEN = 0106    ODD = 1107 108 109class H2ConnectionStateMachine:110    """111    A single HTTP/2 connection state machine.112 113    This state machine, while defined in its own class, is logically part of114    the H2Connection class also defined in this file. The state machine itself115    maintains very little state directly, instead focusing entirely on managing116    state transitions.117    """118 119    # For the purposes of this state machine we treat HEADERS and their120    # associated CONTINUATION frames as a single jumbo frame. The protocol121    # allows/requires this by preventing other frames from being interleved in122    # between HEADERS/CONTINUATION frames.123    #124    # The _transitions dictionary contains a mapping of tuples of125    # (state, input) to tuples of (side_effect_function, end_state). This map126    # contains all allowed transitions: anything not in this map is invalid127    # and immediately causes a transition to ``closed``.128 129    _transitions = {130        # State: idle131        (ConnectionState.IDLE, ConnectionInputs.SEND_HEADERS):132            (None, ConnectionState.CLIENT_OPEN),133        (ConnectionState.IDLE, ConnectionInputs.RECV_HEADERS):134            (None, ConnectionState.SERVER_OPEN),135        (ConnectionState.IDLE, ConnectionInputs.SEND_SETTINGS):136            (None, ConnectionState.IDLE),137        (ConnectionState.IDLE, ConnectionInputs.RECV_SETTINGS):138            (None, ConnectionState.IDLE),139        (ConnectionState.IDLE, ConnectionInputs.SEND_WINDOW_UPDATE):140            (None, ConnectionState.IDLE),141        (ConnectionState.IDLE, ConnectionInputs.RECV_WINDOW_UPDATE):142            (None, ConnectionState.IDLE),143        (ConnectionState.IDLE, ConnectionInputs.SEND_PING):144            (None, ConnectionState.IDLE),145        (ConnectionState.IDLE, ConnectionInputs.RECV_PING):146            (None, ConnectionState.IDLE),147        (ConnectionState.IDLE, ConnectionInputs.SEND_GOAWAY):148            (None, ConnectionState.CLOSED),149        (ConnectionState.IDLE, ConnectionInputs.RECV_GOAWAY):150            (None, ConnectionState.CLOSED),151        (ConnectionState.IDLE, ConnectionInputs.SEND_PRIORITY):152            (None, ConnectionState.IDLE),153        (ConnectionState.IDLE, ConnectionInputs.RECV_PRIORITY):154            (None, ConnectionState.IDLE),155        (ConnectionState.IDLE, ConnectionInputs.SEND_ALTERNATIVE_SERVICE):156            (None, ConnectionState.SERVER_OPEN),157        (ConnectionState.IDLE, ConnectionInputs.RECV_ALTERNATIVE_SERVICE):158            (None, ConnectionState.CLIENT_OPEN),159 160        # State: open, client side.161        (ConnectionState.CLIENT_OPEN, ConnectionInputs.SEND_HEADERS):162            (None, ConnectionState.CLIENT_OPEN),163        (ConnectionState.CLIENT_OPEN, ConnectionInputs.SEND_DATA):164            (None, ConnectionState.CLIENT_OPEN),165        (ConnectionState.CLIENT_OPEN, ConnectionInputs.SEND_GOAWAY):166            (None, ConnectionState.CLOSED),167        (ConnectionState.CLIENT_OPEN, ConnectionInputs.SEND_WINDOW_UPDATE):168            (None, ConnectionState.CLIENT_OPEN),169        (ConnectionState.CLIENT_OPEN, ConnectionInputs.SEND_PING):170            (None, ConnectionState.CLIENT_OPEN),171        (ConnectionState.CLIENT_OPEN, ConnectionInputs.SEND_SETTINGS):172            (None, ConnectionState.CLIENT_OPEN),173        (ConnectionState.CLIENT_OPEN, ConnectionInputs.SEND_PRIORITY):174            (None, ConnectionState.CLIENT_OPEN),175        (ConnectionState.CLIENT_OPEN, ConnectionInputs.RECV_HEADERS):176            (None, ConnectionState.CLIENT_OPEN),177        (ConnectionState.CLIENT_OPEN, ConnectionInputs.RECV_PUSH_PROMISE):178            (None, ConnectionState.CLIENT_OPEN),179        (ConnectionState.CLIENT_OPEN, ConnectionInputs.RECV_DATA):180            (None, ConnectionState.CLIENT_OPEN),181        (ConnectionState.CLIENT_OPEN, ConnectionInputs.RECV_GOAWAY):182            (None, ConnectionState.CLOSED),183        (ConnectionState.CLIENT_OPEN, ConnectionInputs.RECV_WINDOW_UPDATE):184            (None, ConnectionState.CLIENT_OPEN),185        (ConnectionState.CLIENT_OPEN, ConnectionInputs.RECV_PING):186            (None, ConnectionState.CLIENT_OPEN),187        (ConnectionState.CLIENT_OPEN, ConnectionInputs.RECV_SETTINGS):188            (None, ConnectionState.CLIENT_OPEN),189        (ConnectionState.CLIENT_OPEN, ConnectionInputs.SEND_RST_STREAM):190            (None, ConnectionState.CLIENT_OPEN),191        (ConnectionState.CLIENT_OPEN, ConnectionInputs.RECV_RST_STREAM):192            (None, ConnectionState.CLIENT_OPEN),193        (ConnectionState.CLIENT_OPEN, ConnectionInputs.RECV_PRIORITY):194            (None, ConnectionState.CLIENT_OPEN),195        (ConnectionState.CLIENT_OPEN,196            ConnectionInputs.RECV_ALTERNATIVE_SERVICE):197                (None, ConnectionState.CLIENT_OPEN),198 199        # State: open, server side.200        (ConnectionState.SERVER_OPEN, ConnectionInputs.SEND_HEADERS):201            (None, ConnectionState.SERVER_OPEN),202        (ConnectionState.SERVER_OPEN, ConnectionInputs.SEND_PUSH_PROMISE):203            (None, ConnectionState.SERVER_OPEN),204        (ConnectionState.SERVER_OPEN, ConnectionInputs.SEND_DATA):205            (None, ConnectionState.SERVER_OPEN),206        (ConnectionState.SERVER_OPEN, ConnectionInputs.SEND_GOAWAY):207            (None, ConnectionState.CLOSED),208        (ConnectionState.SERVER_OPEN, ConnectionInputs.SEND_WINDOW_UPDATE):209            (None, ConnectionState.SERVER_OPEN),210        (ConnectionState.SERVER_OPEN, ConnectionInputs.SEND_PING):211            (None, ConnectionState.SERVER_OPEN),212        (ConnectionState.SERVER_OPEN, ConnectionInputs.SEND_SETTINGS):213            (None, ConnectionState.SERVER_OPEN),214        (ConnectionState.SERVER_OPEN, ConnectionInputs.SEND_PRIORITY):215            (None, ConnectionState.SERVER_OPEN),216        (ConnectionState.SERVER_OPEN, ConnectionInputs.RECV_HEADERS):217            (None, ConnectionState.SERVER_OPEN),218        (ConnectionState.SERVER_OPEN, ConnectionInputs.RECV_DATA):219            (None, ConnectionState.SERVER_OPEN),220        (ConnectionState.SERVER_OPEN, ConnectionInputs.RECV_GOAWAY):221            (None, ConnectionState.CLOSED),222        (ConnectionState.SERVER_OPEN, ConnectionInputs.RECV_WINDOW_UPDATE):223            (None, ConnectionState.SERVER_OPEN),224        (ConnectionState.SERVER_OPEN, ConnectionInputs.RECV_PING):225            (None, ConnectionState.SERVER_OPEN),226        (ConnectionState.SERVER_OPEN, ConnectionInputs.RECV_SETTINGS):227            (None, ConnectionState.SERVER_OPEN),228        (ConnectionState.SERVER_OPEN, ConnectionInputs.RECV_PRIORITY):229            (None, ConnectionState.SERVER_OPEN),230        (ConnectionState.SERVER_OPEN, ConnectionInputs.SEND_RST_STREAM):231            (None, ConnectionState.SERVER_OPEN),232        (ConnectionState.SERVER_OPEN, ConnectionInputs.RECV_RST_STREAM):233            (None, ConnectionState.SERVER_OPEN),234        (ConnectionState.SERVER_OPEN,235            ConnectionInputs.SEND_ALTERNATIVE_SERVICE):236                (None, ConnectionState.SERVER_OPEN),237        (ConnectionState.SERVER_OPEN,238            ConnectionInputs.RECV_ALTERNATIVE_SERVICE):239                (None, ConnectionState.SERVER_OPEN),240 241        # State: closed242        (ConnectionState.CLOSED, ConnectionInputs.SEND_GOAWAY):243            (None, ConnectionState.CLOSED),244        (ConnectionState.CLOSED, ConnectionInputs.RECV_GOAWAY):245            (None, ConnectionState.CLOSED),246    }247 248    def __init__(self) -> None:249        self.state = ConnectionState.IDLE250 251    def process_input(self, input_: ConnectionInputs) -> list[Event]:252        """253        Process a specific input in the state machine.254        """255        if not isinstance(input_, ConnectionInputs):256            msg = "Input must be an instance of ConnectionInputs"257            raise ValueError(msg)  # noqa: TRY004258 259        try:260            func, target_state = self._transitions[(self.state, input_)]261        except KeyError as e:262            old_state = self.state263            self.state = ConnectionState.CLOSED264            msg = f"Invalid input {input_} in state {old_state}"265            raise ProtocolError(msg) from e266        else:267            self.state = target_state268            if func is not None:  # pragma: no cover269                return func()270 271            return []272 273 274class H2Connection:275    """276    A low-level HTTP/2 connection object. This handles building and receiving277    frames and maintains both connection and per-stream state for all streams278    on this connection.279 280    This wraps a HTTP/2 Connection state machine implementation, ensuring that281    frames can only be sent/received when the connection is in a valid state.282    It also builds stream state machines on demand to ensure that the283    constraints of those state machines are met as well. Attempts to create284    frames that cannot be sent will raise a ``ProtocolError``.285 286    .. versionchanged:: 2.3.0287       Added the ``header_encoding`` keyword argument.288 289    .. versionchanged:: 2.5.0290       Added the ``config`` keyword argument. Deprecated the ``client_side``291       and ``header_encoding`` parameters.292 293    .. versionchanged:: 3.0.0294       Removed deprecated parameters and properties.295 296    :param config: The configuration for the HTTP/2 connection.297 298        .. versionadded:: 2.5.0299 300    :type config: :class:`H2Configuration <h2.config.H2Configuration>`301    """302 303    # The initial maximum outbound frame size. This can be changed by receiving304    # a settings frame.305    DEFAULT_MAX_OUTBOUND_FRAME_SIZE = 65535306 307    # The initial maximum inbound frame size. This is somewhat arbitrarily308    # chosen.309    DEFAULT_MAX_INBOUND_FRAME_SIZE = 2**24310 311    # The highest acceptable stream ID.312    HIGHEST_ALLOWED_STREAM_ID = 2**31 - 1313 314    # The largest acceptable window increment.315    MAX_WINDOW_INCREMENT = 2**31 - 1316 317    # The initial default value of SETTINGS_MAX_HEADER_LIST_SIZE.318    DEFAULT_MAX_HEADER_LIST_SIZE = 2**16319 320    # Keep in memory limited amount of results for streams closes321    MAX_CLOSED_STREAMS = 2**16322 323    def __init__(self, config: H2Configuration | None = None) -> None:324        self.state_machine = H2ConnectionStateMachine()325        self.streams: dict[int, H2Stream] = {}326        self.highest_inbound_stream_id = 0327        self.highest_outbound_stream_id = 0328        self.encoder = Encoder()329        self.decoder = Decoder()330 331        # This won't always actually do anything: for versions of HPACK older332        # than 2.3.0 it does nothing. However, we have to try!333        self.decoder.max_header_list_size = self.DEFAULT_MAX_HEADER_LIST_SIZE334 335        #: The configuration for this HTTP/2 connection object.336        #:337        #: .. versionadded:: 2.5.0338        self.config = config or H2Configuration(client_side=True)339 340        # Objects that store settings, including defaults.341        #342        # We set the MAX_CONCURRENT_STREAMS value to 100 because its default is343        # unbounded, and that's a dangerous default because it allows344        # essentially unbounded resources to be allocated regardless of how345        # they will be used. 100 should be suitable for the average346        # application. This default obviously does not apply to the remote347        # peer's settings: the remote peer controls them!348        #349        # We also set MAX_HEADER_LIST_SIZE to a reasonable value. This is to350        # advertise our defence against CVE-2016-6581. However, not all351        # versions of HPACK will let us do it. That's ok: we should at least352        # suggest that we're not vulnerable.353        self.local_settings = Settings(354            client=self.config.client_side,355            initial_values={356                SettingCodes.MAX_CONCURRENT_STREAMS: 100,357                SettingCodes.MAX_HEADER_LIST_SIZE:358                    self.DEFAULT_MAX_HEADER_LIST_SIZE,359            },360        )361        self.remote_settings = Settings(client=not self.config.client_side)362 363        # The current value of the connection flow control windows on the364        # connection.365        self.outbound_flow_control_window = (366            self.remote_settings.initial_window_size367        )368 369        #: The maximum size of a frame that can be emitted by this peer, in370        #: bytes.371        self.max_outbound_frame_size = self.remote_settings.max_frame_size372 373        #: The maximum size of a frame that can be received by this peer, in374        #: bytes.375        self.max_inbound_frame_size = self.local_settings.max_frame_size376 377        # Buffer for incoming data.378        self.incoming_buffer = FrameBuffer(server=not self.config.client_side)379 380        # A private variable to store a sequence of received header frames381        # until completion.382        self._header_frames: list[Frame] = []383 384        # Data that needs to be sent.385        self._data_to_send = bytearray()386 387        # Keeps track of how streams are closed.388        # Used to ensure that we don't blow up in the face of frames that were389        # in flight when a RST_STREAM was sent.390        # Also used to determine whether we should consider a frame received391        # while a stream is closed as either a stream error or a connection392        # error.393        self._closed_streams: dict[int, StreamClosedBy | None] = SizeLimitDict(394            size_limit=self.MAX_CLOSED_STREAMS,395        )396 397        # The flow control window manager for the connection.398        self._inbound_flow_control_window_manager = WindowManager(399            max_window_size=self.local_settings.initial_window_size,400        )401 402        # When in doubt use dict-dispatch.403        self._frame_dispatch_table: dict[type[Frame], Callable] = {  # type: ignore404            HeadersFrame: self._receive_headers_frame,405            PushPromiseFrame: self._receive_push_promise_frame,406            SettingsFrame: self._receive_settings_frame,407            DataFrame: self._receive_data_frame,408            WindowUpdateFrame: self._receive_window_update_frame,409            PingFrame: self._receive_ping_frame,410            RstStreamFrame: self._receive_rst_stream_frame,411            PriorityFrame: self._receive_priority_frame,412            GoAwayFrame: self._receive_goaway_frame,413            ContinuationFrame: self._receive_naked_continuation,414            AltSvcFrame: self._receive_alt_svc_frame,415            ExtensionFrame: self._receive_unknown_frame,416        }417 418    def _prepare_for_sending(self, frames: list[Frame]) -> None:419        if not frames:420            return421        self._data_to_send += b"".join(f.serialize() for f in frames)422        assert all(f.body_len <= self.max_outbound_frame_size for f in frames)423 424    def _open_streams(self, remainder: int) -> int:425        """426        A common method of counting number of open streams. Returns the number427        of streams that are open *and* that have (stream ID % 2) == remainder.428        While it iterates, also deletes any closed streams.429        """430        count = 0431        to_delete = []432 433        for stream_id, stream in self.streams.items():434            if stream.open and (stream_id % 2 == remainder):435                count += 1436            elif stream.closed:437                to_delete.append(stream_id)438 439        for stream_id in to_delete:440            stream = self.streams.pop(stream_id)441            self._closed_streams[stream_id] = stream.closed_by442 443        return count444 445    @property446    def open_outbound_streams(self) -> int:447        """448        The current number of open outbound streams.449        """450        outbound_numbers = int(self.config.client_side)451        return self._open_streams(outbound_numbers)452 453    @property454    def open_inbound_streams(self) -> int:455        """456        The current number of open inbound streams.457        """458        inbound_numbers = int(not self.config.client_side)459        return self._open_streams(inbound_numbers)460 461    @property462    def inbound_flow_control_window(self) -> int:463        """464        The size of the inbound flow control window for the connection. This is465        rarely publicly useful: instead, use :meth:`remote_flow_control_window466        <h2.connection.H2Connection.remote_flow_control_window>`. This467        shortcut is largely present to provide a shortcut to this data.468        """469        return self._inbound_flow_control_window_manager.current_window_size470 471    def _begin_new_stream(self, stream_id: int, allowed_ids: AllowedStreamIDs) -> H2Stream:472        """473        Initiate a new stream.474 475        .. versionchanged:: 2.0.0476           Removed this function from the public API.477 478        :param stream_id: The ID of the stream to open.479        :param allowed_ids: What kind of stream ID is allowed.480        """481        self.config.logger.debug(482            "Attempting to initiate stream ID %d", stream_id,483        )484        outbound = self._stream_id_is_outbound(stream_id)485        highest_stream_id = (486            self.highest_outbound_stream_id if outbound else487            self.highest_inbound_stream_id488        )489 490        if stream_id <= highest_stream_id:491            raise StreamIDTooLowError(stream_id, highest_stream_id)492 493        if (stream_id % 2) != int(allowed_ids):494            msg = "Invalid stream ID for peer."495            raise ProtocolError(msg)496 497        s = H2Stream(498            stream_id,499            config=self.config,500            inbound_window_size=self.local_settings.initial_window_size,501            outbound_window_size=self.remote_settings.initial_window_size,502        )503        self.config.logger.debug("Stream ID %d created", stream_id)504        s.max_outbound_frame_size = self.max_outbound_frame_size505 506        self.streams[stream_id] = s507        self.config.logger.debug("Current streams: %s", self.streams.keys())508 509        if outbound:510            self.highest_outbound_stream_id = stream_id511        else:512            self.highest_inbound_stream_id = stream_id513 514        return s515 516    def initiate_connection(self) -> None:517        """518        Provides any data that needs to be sent at the start of the connection.519        Must be called for both clients and servers.520        """521        self.config.logger.debug("Initializing connection")522        self.state_machine.process_input(ConnectionInputs.SEND_SETTINGS)523        if self.config.client_side:524            preamble = b"PRI * HTTP/2.0\r\n\r\nSM\r\n\r\n"525        else:526            preamble = b""527 528        f = SettingsFrame(0)529        for setting, value in self.local_settings.items():530            f.settings[setting] = value531        self.config.logger.debug(532            "Send Settings frame: %s", self.local_settings,533        )534 535        self._data_to_send += preamble + f.serialize()536 537    def initiate_upgrade_connection(self, settings_header: bytes | None = None) -> bytes | None:538        """539        Call to initialise the connection object for use with an upgraded540        HTTP/2 connection (i.e. a connection negotiated using the541        ``Upgrade: h2c`` HTTP header).542 543        This method differs from :meth:`initiate_connection544        <h2.connection.H2Connection.initiate_connection>` in several ways.545        Firstly, it handles the additional SETTINGS frame that is sent in the546        ``HTTP2-Settings`` header field. When called on a client connection,547        this method will return a bytestring that the caller can put in the548        ``HTTP2-Settings`` field they send on their initial request. When549        called on a server connection, the user **must** provide the value they550        received from the client in the ``HTTP2-Settings`` header field to the551        ``settings_header`` argument, which will be used appropriately.552 553        Additionally, this method sets up stream 1 in a half-closed state554        appropriate for this side of the connection, to reflect the fact that555        the request is already complete.556 557        Finally, this method also prepares the appropriate preamble to be sent558        after the upgrade.559 560        .. versionadded:: 2.3.0561 562        :param settings_header: (optional, server-only): The value of the563             ``HTTP2-Settings`` header field received from the client.564        :type settings_header: ``bytes``565 566        :returns: For clients, a bytestring to put in the ``HTTP2-Settings``.567            For servers, returns nothing.568        :rtype: ``bytes`` or ``None``569        """570        self.config.logger.debug(571            "Upgrade connection. Current settings: %s", self.local_settings,572        )573 574        frame_data = None575        # Begin by getting the preamble in place.576        self.initiate_connection()577 578        if self.config.client_side:579            f = SettingsFrame(0)580            for setting, value in self.local_settings.items():581                f.settings[setting] = value582 583            frame_data = f.serialize_body()584            frame_data = base64.urlsafe_b64encode(frame_data)585        elif settings_header:586            # We have a settings header from the client. This needs to be587            # applied, but we want to throw away the ACK. We do this by588            # inserting the data into a Settings frame and then passing it to589            # the state machine, but ignoring the return value.590            settings_header = base64.urlsafe_b64decode(settings_header)591            f = SettingsFrame(0)592            f.parse_body(memoryview(settings_header))593            self._receive_settings_frame(f)594 595        # Set up appropriate state. Stream 1 in a half-closed state:596        # half-closed(local) for clients, half-closed(remote) for servers.597        # Additionally, we need to set up the Connection state machine.598        connection_input = (599            ConnectionInputs.SEND_HEADERS if self.config.client_side600            else ConnectionInputs.RECV_HEADERS601        )602        self.config.logger.debug("Process input %s", connection_input)603        self.state_machine.process_input(connection_input)604 605        # Set up stream 1.606        self._begin_new_stream(stream_id=1, allowed_ids=AllowedStreamIDs.ODD)607        self.streams[1].upgrade(self.config.client_side)608        return frame_data609 610    def _get_or_create_stream(self, stream_id: int, allowed_ids: AllowedStreamIDs) -> H2Stream:611        """612        Gets a stream by its stream ID. Will create one if one does not already613        exist. Use allowed_ids to circumvent the usual stream ID rules for614        clients and servers.615 616        .. versionchanged:: 2.0.0617           Removed this function from the public API.618        """619        try:620            return self.streams[stream_id]621        except KeyError:622            return self._begin_new_stream(stream_id, allowed_ids)623 624    def _get_stream_by_id(self, stream_id: int | None) -> H2Stream:625        """626        Gets a stream by its stream ID. Raises NoSuchStreamError if the stream627        ID does not correspond to a known stream and is higher than the current628        maximum: raises if it is lower than the current maximum.629 630        .. versionchanged:: 2.0.0631           Removed this function from the public API.632        """633        if not stream_id:634            raise NoSuchStreamError(-1)  # pragma: no cover635        try:636            return self.streams[stream_id]637        except KeyError as e:638            outbound = self._stream_id_is_outbound(stream_id)639            highest_stream_id = (640                self.highest_outbound_stream_id if outbound else641                self.highest_inbound_stream_id642            )643 644            if stream_id > highest_stream_id:645                raise NoSuchStreamError(stream_id) from e646            raise StreamClosedError(stream_id) from e647 648    def get_next_available_stream_id(self) -> int:649        """650        Returns an integer suitable for use as the stream ID for the next651        stream created by this endpoint. For server endpoints, this stream ID652        will be even. For client endpoints, this stream ID will be odd. If no653        stream IDs are available, raises :class:`NoAvailableStreamIDError654        <h2.exceptions.NoAvailableStreamIDError>`.655 656        .. warning:: The return value from this function does not change until657                     the stream ID has actually been used by sending or pushing658                     headers on that stream. For that reason, it should be659                     called as close as possible to the actual use of the660                     stream ID.661 662        .. versionadded:: 2.0.0663 664        :raises: :class:`NoAvailableStreamIDError665            <h2.exceptions.NoAvailableStreamIDError>`666        :returns: The next free stream ID this peer can use to initiate a667            stream.668        :rtype: ``int``669        """670        # No streams have been opened yet, so return the lowest allowed stream671        # ID.672        if not self.highest_outbound_stream_id:673            next_stream_id = 1 if self.config.client_side else 2674        else:675            next_stream_id = self.highest_outbound_stream_id + 2676        self.config.logger.debug(677            "Next available stream ID %d", next_stream_id,678        )679        if next_stream_id > self.HIGHEST_ALLOWED_STREAM_ID:680            msg = "Exhausted allowed stream IDs"681            raise NoAvailableStreamIDError(msg)682 683        return next_stream_id684 685    def send_headers(self,686                     stream_id: int,687                     headers: Iterable[HeaderWeaklyTyped],688                     end_stream: bool = False,689                     priority_weight: int | None = None,690                     priority_depends_on: int | None = None,691                     priority_exclusive: bool | None = None) -> None:692        """693        Send headers on a given stream.694 695        This function can be used to send request or response headers: the kind696        that are sent depends on whether this connection has been opened as a697        client or server connection, and whether the stream was opened by the698        remote peer or not.699 700        If this is a client connection, calling ``send_headers`` will send the701        headers as a request. It will also implicitly open the stream being702        used. If this is a client connection and ``send_headers`` has *already*703        been called, this will send trailers instead.704 705        If this is a server connection, calling ``send_headers`` will send the706        headers as a response. It is a protocol error for a server to open a707        stream by sending headers. If this is a server connection and708        ``send_headers`` has *already* been called, this will send trailers709        instead.710 711        When acting as a server, you may call ``send_headers`` any number of712        times allowed by the following rules, in this order:713 714        - zero or more times with ``(':status', '1XX')`` (where ``1XX`` is a715          placeholder for any 100-level status code).716        - once with any other status header.717        - zero or one time for trailers.718 719        That is, you are allowed to send as many informational responses as you720        like, followed by one complete response and zero or one HTTP trailer721        blocks.722 723        Clients may send one or two header blocks: one request block, and724        optionally one trailer block.725 726        If it is important to send HPACK "never indexed" header fields (as727        defined in `RFC 7451 Section 7.1.3728        <https://tools.ietf.org/html/rfc7541#section-7.1.3>`_), the user may729        instead provide headers using the HPACK library's :class:`HeaderTuple730        <hpack:hpack.HeaderTuple>` and :class:`NeverIndexedHeaderTuple731        <hpack:hpack.NeverIndexedHeaderTuple>` objects.732 733        This method also allows users to prioritize the stream immediately,734        by sending priority information on the HEADERS frame directly. To do735        this, any one of ``priority_weight``, ``priority_depends_on``, or736        ``priority_exclusive`` must be set to a value that is not ``None``. For737        more information on the priority fields, see :meth:`prioritize738        <h2.connection.H2Connection.prioritize>`.739 740        .. warning:: In HTTP/2, it is mandatory that all the HTTP/2 special741            headers (that is, ones whose header keys begin with ``:``) appear742            at the start of the header block, before any normal headers.743 744        .. versionchanged:: 2.3.0745           Added support for using :class:`HeaderTuple746           <hpack:hpack.HeaderTuple>` objects to store headers.747 748        .. versionchanged:: 2.4.0749           Added the ability to provide priority keyword arguments:750           ``priority_weight``, ``priority_depends_on``, and751           ``priority_exclusive``.752 753        :param stream_id: The stream ID to send the headers on. If this stream754            does not currently exist, it will be created.755        :type stream_id: ``int``756 757        :param headers: The request/response headers to send.758        :type headers: An iterable of two tuples of bytestrings or759            :class:`HeaderTuple <hpack:hpack.HeaderTuple>` objects.760 761        :param end_stream: Whether this headers frame should end the stream762            immediately (that is, whether no more data will be sent after this763            frame). Defaults to ``False``.764        :type end_stream: ``bool``765 766        :param priority_weight: Sets the priority weight of the stream. See767            :meth:`prioritize <h2.connection.H2Connection.prioritize>` for more768            about how this field works. Defaults to ``None``, which means that769            no priority information will be sent.770        :type priority_weight: ``int`` or ``None``771 772        :param priority_depends_on: Sets which stream this one depends on for773            priority purposes. See :meth:`prioritize774            <h2.connection.H2Connection.prioritize>` for more about how this775            field works. Defaults to ``None``, which means that no priority776            information will be sent.777        :type priority_depends_on: ``int`` or ``None``778 779        :param priority_exclusive: Sets whether this stream exclusively depends780            on the stream given in ``priority_depends_on`` for priority781            purposes. See :meth:`prioritize782            <h2.connection.H2Connection.prioritize>` for more about how this783            field workds. Defaults to ``None``, which means that no priority784            information will be sent.785        :type priority_depends_on: ``bool`` or ``None``786 787        :returns: Nothing788        """789        self.config.logger.debug(790            "Send headers on stream ID %d", stream_id,791        )792 793        # Check we can open the stream.794        if stream_id not in self.streams:795            max_open_streams = self.remote_settings.max_concurrent_streams796            value = self.open_outbound_streams # take a copy due to the property accessor having side affects797            if (value + 1) > max_open_streams:798                msg = f"Max outbound streams is {max_open_streams}, {value} open"799                raise TooManyStreamsError(msg)800 801        self.state_machine.process_input(ConnectionInputs.SEND_HEADERS)802        stream = self._get_or_create_stream(803            stream_id, AllowedStreamIDs(self.config.client_side),804        )805 806        frames: list[Frame] = []807        frames.extend(stream.send_headers(808            headers, self.encoder, end_stream,809        ))810 811        # We may need to send priority information.812        priority_present = (813            (priority_weight is not None) or814            (priority_depends_on is not None) or815            (priority_exclusive is not None)816        )817 818        if priority_present:819            if not self.config.client_side:820                msg = "Servers SHOULD NOT prioritize streams."821                raise RFC1122Error(msg)822 823            headers_frame = frames[0]824            assert isinstance(headers_frame, HeadersFrame)825 826            headers_frame.flags.add("PRIORITY")827            frames[0] = _add_frame_priority(828                headers_frame,829                priority_weight,830                priority_depends_on,831                priority_exclusive,832            )833 834        self._prepare_for_sending(frames)835 836    def send_data(self,837                  stream_id: int,838                  data: bytes | memoryview,839                  end_stream: bool = False,840                  pad_length: Any = None) -> None:841        """842        Send data on a given stream.843 844        This method does no breaking up of data: if the data is larger than the845        value returned by :meth:`local_flow_control_window846        <h2.connection.H2Connection.local_flow_control_window>` for this stream847        then a :class:`FlowControlError <h2.exceptions.FlowControlError>` will848        be raised. If the data is larger than :data:`max_outbound_frame_size849        <h2.connection.H2Connection.max_outbound_frame_size>` then a850        :class:`FrameTooLargeError <h2.exceptions.FrameTooLargeError>` will be851        raised.852 853        h2 does this to avoid buffering the data internally. If the user854        has more data to send than h2 will allow, consider breaking it up855        and buffering it externally.856 857        :param stream_id: The ID of the stream on which to send the data.858        :type stream_id: ``int``859        :param data: The data to send on the stream.860        :type data: ``bytes``861        :param end_stream: (optional) Whether this is the last data to be sent862            on the stream. Defaults to ``False``.863        :type end_stream: ``bool``864        :param pad_length: (optional) Length of the padding to apply to the865            data frame. Defaults to ``None`` for no use of padding. Note that866            a value of ``0`` results in padding of length ``0``867            (with the "padding" flag set on the frame).868 869            .. versionadded:: 2.6.0870 871        :type pad_length: ``int``872        :returns: Nothing873        """874        self.config.logger.debug(875            "Send data on stream ID %d with len %d", stream_id, len(data),876        )877        frame_size = len(data)878        if pad_length is not None:879            if not isinstance(pad_length, int):880                msg = "pad_length must be an int"881                raise TypeError(msg)882            if pad_length < 0 or pad_length > 255:883                msg = "pad_length must be within range: [0, 255]"884                raise ValueError(msg)885            # Account for padding bytes plus the 1-byte padding length field.886            frame_size += pad_length + 1887        self.config.logger.debug(888            "Frame size on stream ID %d is %d", stream_id, frame_size,889        )890 891        if frame_size > self.local_flow_control_window(stream_id):892            msg = f"Cannot send {frame_size} bytes, flow control window is {self.local_flow_control_window(stream_id)}"893            raise FlowControlError(msg)894        if frame_size > self.max_outbound_frame_size:895            msg = f"Cannot send frame size {frame_size}, max frame size is {self.max_outbound_frame_size}"896            raise FrameTooLargeError(msg)897 898        self.state_machine.process_input(ConnectionInputs.SEND_DATA)899        frames = self.streams[stream_id].send_data(900            data, end_stream, pad_length=pad_length,901        )902 903        self._prepare_for_sending(frames)904 905        self.outbound_flow_control_window -= frame_size906        self.config.logger.debug(907            "Outbound flow control window size is %d",908            self.outbound_flow_control_window,909        )910        assert self.outbound_flow_control_window >= 0911 912    def end_stream(self, stream_id: int) -> None:913        """914        Cleanly end a given stream.915 916        This method ends a stream by sending an empty DATA frame on that stream917        with the ``END_STREAM`` flag set.918 919        :param stream_id: The ID of the stream to end.920        :type stream_id: ``int``921        :returns: Nothing922        """923        self.config.logger.debug("End stream ID %d", stream_id)924        self.state_machine.process_input(ConnectionInputs.SEND_DATA)925        frames = self.streams[stream_id].end_stream()926        self._prepare_for_sending(frames)927 928    def increment_flow_control_window(self, increment: int, stream_id: int | None = None) -> None:929        """930        Increment a flow control window, optionally for a single stream. Allows931        the remote peer to send more data.932 933        .. versionchanged:: 2.0.0934           Rejects attempts to increment the flow control window by out of935           range values with a ``ValueError``.936 937        :param increment: The amount to increment the flow control window by.938        :type increment: ``int``939        :param stream_id: (optional) The ID of the stream that should have its940            flow control window opened. If not present or ``None``, the941            connection flow control window will be opened instead.942        :type stream_id: ``int`` or ``None``943        :returns: Nothing944        :raises: ``ValueError``945        """946        if not (1 <= increment <= self.MAX_WINDOW_INCREMENT):947            msg = f"Flow control increment must be between 1 and {self.MAX_WINDOW_INCREMENT}"948            raise ValueError(msg)949 950        self.state_machine.process_input(ConnectionInputs.SEND_WINDOW_UPDATE)951 952        if stream_id is not None:953            stream = self.streams[stream_id]954            frames = stream.increase_flow_control_window(955                increment,956            )957 958            self.config.logger.debug(959                "Increase stream ID %d flow control window by %d",960                stream_id, increment,961            )962        else:963            self._inbound_flow_control_window_manager.window_opened(increment)964            f = WindowUpdateFrame(0)965            f.window_increment = increment966            frames = [f]967 968            self.config.logger.debug(969                "Increase connection flow control window by %d", increment,970            )971 972        self._prepare_for_sending(frames)973 974    def push_stream(self,975                    stream_id: int,976                    promised_stream_id: int,977                    request_headers: Iterable[HeaderWeaklyTyped]) -> None:978        """979        Push a response to the client by sending a PUSH_PROMISE frame.980 981        If it is important to send HPACK "never indexed" header fields (as982        defined in `RFC 7451 Section 7.1.3983        <https://tools.ietf.org/html/rfc7541#section-7.1.3>`_), the user may984        instead provide headers using the HPACK library's :class:`HeaderTuple985        <hpack:hpack.HeaderTuple>` and :class:`NeverIndexedHeaderTuple986        <hpack:hpack.NeverIndexedHeaderTuple>` objects.987 988        :param stream_id: The ID of the stream that this push is a response to.989        :type stream_id: ``int``990        :param promised_stream_id: The ID of the stream that the pushed991            response will be sent on.992        :type promised_stream_id: ``int``993        :param request_headers: The headers of the request that the pushed994            response will be responding to.995        :type request_headers: An iterable of two tuples of bytestrings or996            :class:`HeaderTuple <hpack:hpack.HeaderTuple>` objects.997        :returns: Nothing998        """999        self.config.logger.debug(1000            "Send Push Promise frame on stream ID %d", stream_id,1001        )1002 1003        if not self.remote_settings.enable_push:1004            msg = "Remote peer has disabled stream push"1005            raise ProtocolError(msg)1006 1007        self.state_machine.process_input(ConnectionInputs.SEND_PUSH_PROMISE)1008        stream = self._get_stream_by_id(stream_id)1009 1010        # We need to prevent users pushing streams in response to streams that1011        # they themselves have already pushed: see #163 and RFC 7540 § 6.6. The1012        # easiest way to do that is to assert that the stream_id is not even:1013        # this shortcut works because only servers can push and the state1014        # machine will enforce this.1015        if (stream_id % 2) == 0:1016            msg = "Cannot recursively push streams."1017            raise ProtocolError(msg)1018 1019        new_stream = self._begin_new_stream(1020            promised_stream_id, AllowedStreamIDs.EVEN,1021        )1022        self.streams[promised_stream_id] = new_stream1023 1024        frames = stream.push_stream_in_band(1025            promised_stream_id, request_headers, self.encoder,1026        )1027        new_frames = new_stream.locally_pushed()1028        self._prepare_for_sending(frames + new_frames)1029 1030    def ping(self, opaque_data: bytes | str) -> None:1031        """1032        Send a PING frame.1033 1034        :param opaque_data: A bytestring of length 8 that will be sent in the1035                            PING frame.1036        :returns: Nothing1037        """1038        self.config.logger.debug("Send Ping frame")1039 1040        if not isinstance(opaque_data, bytes) or len(opaque_data) != 8:1041            msg = f"Invalid value for ping data: {opaque_data!r}"1042            raise ValueError(msg)1043 1044        self.state_machine.process_input(ConnectionInputs.SEND_PING)1045        f = PingFrame(0)1046        f.opaque_data = opaque_data1047        self._prepare_for_sending([f])1048 1049    def reset_stream(self, stream_id: int, error_code: ErrorCodes | int = 0) -> None:1050        """1051        Reset a stream.1052 1053        This method forcibly closes a stream by sending a RST_STREAM frame for1054        a given stream. This is not a graceful closure. To gracefully end a1055        stream, try the :meth:`end_stream1056        <h2.connection.H2Connection.end_stream>` method.1057 1058        :param stream_id: The ID of the stream to reset.1059        :type stream_id: ``int``1060        :param error_code: (optional) The error code to use to reset the1061            stream. Defaults to :data:`ErrorCodes.NO_ERROR1062            <h2.errors.ErrorCodes.NO_ERROR>`.1063        :type error_code: ``int``1064        :returns: Nothing1065        """1066        self.config.logger.debug("Reset stream ID %d", stream_id)1067        self.state_machine.process_input(ConnectionInputs.SEND_RST_STREAM)1068        stream = self._get_stream_by_id(stream_id)1069        frames = stream.reset_stream(error_code)1070        self._prepare_for_sending(frames)1071 1072    def close_connection(self,1073                         error_code: ErrorCodes | int = 0,1074                         additional_data: bytes | None = None,1075                         last_stream_id: int | None = None) -> None:1076        """1077        Close a connection, emitting a GOAWAY frame.1078 1079        .. versionchanged:: 2.4.01080           Added ``additional_data`` and ``last_stream_id`` arguments.1081 1082        :param error_code: (optional) The error code to send in the GOAWAY1083            frame.1084        :param additional_data: (optional) Additional debug data indicating1085            a reason for closing the connection. Must be a bytestring.1086        :param last_stream_id: (optional) The last stream which was processed1087            by the sender. Defaults to ``highest_inbound_stream_id``.1088        :returns: Nothing1089        """1090        self.config.logger.debug("Close connection")1091        self.state_machine.process_input(ConnectionInputs.SEND_GOAWAY)1092 1093        # Additional_data must be bytes1094        if additional_data is not None:1095            assert isinstance(additional_data, bytes)1096 1097        if last_stream_id is None:1098            last_stream_id = self.highest_inbound_stream_id1099 1100        f = GoAwayFrame(1101            stream_id=0,1102            last_stream_id=last_stream_id,1103            error_code=error_code,1104            additional_data=(additional_data or b""),1105        )1106        self._prepare_for_sending([f])1107 1108    def update_settings(self, new_settings: dict[SettingCodes | int, int]) -> None:1109        """1110        Update the local settings. This will prepare and emit the appropriate1111        SETTINGS frame.1112 1113        :param new_settings: A dictionary of {setting: new value}1114        """1115        self.config.logger.debug(1116            "Update connection settings to %s", new_settings,1117        )1118        self.state_machine.process_input(ConnectionInputs.SEND_SETTINGS)1119        self.local_settings.update(new_settings)1120        s = SettingsFrame(0)1121        s.settings = new_settings1122        self._prepare_for_sending([s])1123 1124    def advertise_alternative_service(self,1125                                      field_value: bytes | str,1126                                      origin: bytes | None = None,1127                                      stream_id: int | None = None) -> None:1128        """1129        Notify a client about an available Alternative Service.1130 1131        An Alternative Service is defined in `RFC 78381132        <https://tools.ietf.org/html/rfc7838>`_. An Alternative Service1133        notification informs a client that a given origin is also available1134        elsewhere.1135 1136        Alternative Services can be advertised in two ways. Firstly, they can1137        be advertised explicitly: that is, a server can say "origin X is also1138        available at Y". To advertise like this, set the ``origin`` argument1139        and not the ``stream_id`` argument. Alternatively, they can be1140        advertised implicitly: that is, a server can say "the origin you're1141        contacting on stream X is also available at Y". To advertise like this,1142        set the ``stream_id`` argument and not the ``origin`` argument.1143 1144        The explicit method of advertising can be done as long as the1145        connection is active. The implicit method can only be done after the1146        client has sent the request headers and before the server has sent the1147        response headers: outside of those points, h2 will forbid sending1148        the Alternative Service advertisement by raising a ProtocolError.1149 1150        The ``field_value`` parameter is specified in RFC 7838. h2 does1151        not validate or introspect this argument: the user is required to1152        ensure that it's well-formed. ``field_value`` corresponds to RFC 7838's1153        "Alternative Service Field Value".1154 1155        .. note:: It is strongly preferred to use the explicit method of1156                  advertising Alternative Services. The implicit method of1157                  advertising Alternative Services has a number of subtleties1158                  and can lead to inconsistencies between the server and1159                  client. h2 allows both mechanisms, but caution is1160                  strongly advised.1161 1162        .. versionadded:: 2.3.01163 1164        :param field_value: The RFC 7838 Alternative Service Field Value. This1165            argument is not introspected by h2: the user is responsible1166            for ensuring that it is well-formed.1167        :type field_value: ``bytes``1168 1169        :param origin: The origin/authority to which the Alternative Service1170            being advertised applies. Must not be provided at the same time as1171            ``stream_id``.1172        :type origin: ``bytes`` or ``None``1173 1174        :param stream_id: The ID of the stream which was sent to the authority1175            for which this Alternative Service advertisement applies. Must not1176            be provided at the same time as ``origin``.1177        :type stream_id: ``int`` or ``None``1178 1179        :returns: Nothing.1180        """1181        if not isinstance(field_value, bytes):1182            msg = "Field must be bytestring."1183            raise ValueError(msg)  # noqa: TRY0041184 1185        if origin is not None and stream_id is not None:1186            msg = "Must not provide both origin and stream_id"1187            raise ValueError(msg)1188 1189        self.state_machine.process_input(1190            ConnectionInputs.SEND_ALTERNATIVE_SERVICE,1191        )1192 1193        if origin is not None:1194            # This ALTSVC is sent on stream zero.1195            f = AltSvcFrame(stream_id=0)1196            f.origin = origin1197            f.field = field_value1198            frames: list[Frame] = [f]1199        else:1200            stream = self._get_stream_by_id(stream_id)

Showing the first 1,200 of 2111 lines. Download the file for the rest.