codekingpro/portable-devtools
114k
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)