Team Ai
Datasetpublic

codekingpro/portable-devtools

sourceHugging Faceupdated 5mo agoView on Hugging Face
1likes14kdownloads
stream.py1426 linesDownload Raw Back to h2
1"""2h2/stream3~~~~~~~~~4 5An implementation of a HTTP/2 stream.6"""7from __future__ import annotations8 9from enum import Enum, IntEnum10from typing import TYPE_CHECKING, Any, Union, cast11 12from hpack import HeaderTuple13from hyperframe.frame import AltSvcFrame, ContinuationFrame, DataFrame, Frame, HeadersFrame, PushPromiseFrame, RstStreamFrame, WindowUpdateFrame14 15from .errors import ErrorCodes, _error_code_from_int16from .events import (17    AlternativeServiceAvailable,18    DataReceived,19    Event,20    InformationalResponseReceived,21    PushedStreamReceived,22    RequestReceived,23    ResponseReceived,24    StreamEnded,25    StreamReset,26    TrailersReceived,27    WindowUpdated,28    _PushedRequestSent,29    _RequestSent,30    _ResponseSent,31    _TrailersSent,32)33from .exceptions import FlowControlError, InvalidBodyLengthError, ProtocolError, StreamClosedError34from .utilities import (35    HeaderValidationFlags,36    authority_from_headers,37    extract_method_header,38    guard_increment_window,39    is_informational_response,40    normalize_inbound_headers,41    normalize_outbound_headers,42    utf8_encode_headers,43    validate_headers,44    validate_outbound_headers,45)46from .windows import WindowManager47 48if TYPE_CHECKING:  # pragma: no cover49    from collections.abc import Callable, Generator, Iterable50 51    from hpack.hpack import Encoder52    from hpack.struct import Header, HeaderWeaklyTyped53 54    from .config import H2Configuration55 56 57class StreamState(IntEnum):58    IDLE = 059    RESERVED_REMOTE = 160    RESERVED_LOCAL = 261    OPEN = 362    HALF_CLOSED_REMOTE = 463    HALF_CLOSED_LOCAL = 564    CLOSED = 665 66 67class StreamInputs(Enum):68    SEND_HEADERS = 069    SEND_PUSH_PROMISE = 170    SEND_RST_STREAM = 271    SEND_DATA = 372    SEND_WINDOW_UPDATE = 473    SEND_END_STREAM = 574    RECV_HEADERS = 675    RECV_PUSH_PROMISE = 776    RECV_RST_STREAM = 877    RECV_DATA = 978    RECV_WINDOW_UPDATE = 1079    RECV_END_STREAM = 1180    RECV_CONTINUATION = 12  # Added in 2.0.081    SEND_INFORMATIONAL_HEADERS = 13  # Added in 2.2.082    RECV_INFORMATIONAL_HEADERS = 14  # Added in 2.2.083    SEND_ALTERNATIVE_SERVICE = 15  # Added in 2.3.084    RECV_ALTERNATIVE_SERVICE = 16  # Added in 2.3.085    UPGRADE_CLIENT = 17  # Added 2.3.086    UPGRADE_SERVER = 18  # Added 2.3.087 88 89class StreamClosedBy(Enum):90    SEND_END_STREAM = 091    RECV_END_STREAM = 192    SEND_RST_STREAM = 293    RECV_RST_STREAM = 394 95 96# This array is initialized once, and is indexed by the stream states above.97# It indicates whether a stream in the given state is open. The reason we do98# this is that we potentially check whether a stream in a given state is open99# quite frequently: given that we check so often, we should do so in the100# fastest and most performant way possible.101STREAM_OPEN = [False for _ in range(len(StreamState))]102STREAM_OPEN[StreamState.OPEN] = True103STREAM_OPEN[StreamState.HALF_CLOSED_LOCAL] = True104STREAM_OPEN[StreamState.HALF_CLOSED_REMOTE] = True105 106 107class H2StreamStateMachine:108    """109    A single HTTP/2 stream state machine.110 111    This stream object implements basically the state machine described in112    RFC 7540 section 5.1.113 114    :param stream_id: The stream ID of this stream. This is stored primarily115        for logging purposes.116    """117 118    def __init__(self, stream_id: int) -> None:119        self.state = StreamState.IDLE120        self.stream_id = stream_id121 122        #: Whether this peer is the client side of this stream.123        self.client: bool | None = None124 125        # Whether trailers have been sent/received on this stream or not.126        self.headers_sent: bool | None = None127        self.trailers_sent: bool | None = None128        self.headers_received: bool | None = None129        self.trailers_received: bool | None = None130 131        # How the stream was closed. One of StreamClosedBy.132        self.stream_closed_by: StreamClosedBy | None = None133 134    def process_input(self, input_: StreamInputs) -> list[Event]:135        """136        Process a specific input in the state machine.137        """138        if not isinstance(input_, StreamInputs):139            msg = "Input must be an instance of StreamInputs"140            raise ValueError(msg)  # noqa: TRY004141 142        try:143            func, target_state = _transitions[(self.state, input_)]144        except KeyError as err:145            old_state = self.state146            self.state = StreamState.CLOSED147            msg = f"Invalid input {input_} in state {old_state}"148            raise ProtocolError(msg) from err149        else:150            previous_state = self.state151            self.state = target_state152            if func is not None:153                try:154                    return func(self, previous_state)155                except ProtocolError:156                    self.state = StreamState.CLOSED157                    raise158                except AssertionError as err:  # pragma: no cover159                    self.state = StreamState.CLOSED160                    raise ProtocolError(err) from err161 162            return []163 164    def request_sent(self, previous_state: StreamState) -> list[Event]:165        """166        Fires when a request is sent.167        """168        self.client = True169        self.headers_sent = True170        event = _RequestSent()171 172        return [event]173 174    def response_sent(self, previous_state: StreamState) -> list[Event]:175        """176        Fires when something that should be a response is sent. This 'response'177        may actually be trailers.178        """179        if not self.headers_sent:180            if self.client is True or self.client is None:181                msg = "Client cannot send responses."182                raise ProtocolError(msg)183            self.headers_sent = True184            return [_ResponseSent()]185        assert not self.trailers_sent186        self.trailers_sent = True187        return [_TrailersSent()]188 189    def request_received(self, previous_state: StreamState) -> list[Event]:190        """191        Fires when a request is received.192        """193        assert not self.headers_received194        assert not self.trailers_received195 196        self.client = False197        self.headers_received = True198        event = RequestReceived(stream_id=self.stream_id)199        return [event]200 201    def response_received(self, previous_state: StreamState) -> list[Event]:202        """203        Fires when a response is received. Also disambiguates between responses204        and trailers.205        """206        event: ResponseReceived | TrailersReceived207        if not self.headers_received:208            assert self.client is True209            self.headers_received = True210            event = ResponseReceived(stream_id=self.stream_id)211        else:212            assert not self.trailers_received213            self.trailers_received = True214            event = TrailersReceived(stream_id=self.stream_id)215 216        event.stream_id = self.stream_id217        return [event]218 219    def data_received(self, previous_state: StreamState) -> list[Event]:220        """221        Fires when data is received.222        """223        if not self.headers_received:224            msg = "cannot receive data before headers"225            raise ProtocolError(msg)226        event = DataReceived(stream_id=self.stream_id)227        return [event]228 229    def window_updated(self, previous_state: StreamState) -> list[Event]:230        """231        Fires when a window update frame is received.232        """233        return [WindowUpdated(stream_id=self.stream_id)]234 235    def stream_half_closed(self, previous_state: StreamState) -> list[Event]:236        """237        Fires when an END_STREAM flag is received in the OPEN state,238        transitioning this stream to a HALF_CLOSED_REMOTE state.239        """240        event = StreamEnded(stream_id=self.stream_id)241        return [event]242 243    def stream_ended(self, previous_state: StreamState) -> list[Event]:244        """245        Fires when a stream is cleanly ended.246        """247        self.stream_closed_by = StreamClosedBy.RECV_END_STREAM248        event = StreamEnded(stream_id=self.stream_id)249        return [event]250 251    def stream_reset(self, previous_state: StreamState) -> list[Event]:252        """253        Fired when a stream is forcefully reset.254        """255        self.stream_closed_by = StreamClosedBy.RECV_RST_STREAM256        return [StreamReset(stream_id=self.stream_id)]257 258    def send_new_pushed_stream(self, previous_state: StreamState) -> list[Event]:259        """260        Fires on the newly pushed stream, when pushed by the local peer.261 262        No event here, but definitionally this peer must be a server.263        """264        assert self.client is None265        self.client = False266        self.headers_received = True267        return []268 269    def recv_new_pushed_stream(self, previous_state: StreamState) -> list[Event]:270        """271        Fires on the newly pushed stream, when pushed by the remote peer.272 273        No event here, but definitionally this peer must be a client.274        """275        assert self.client is None276        self.client = True277        self.headers_sent = True278        return []279 280    def send_push_promise(self, previous_state: StreamState) -> list[Event]:281        """282        Fires on the already-existing stream when a PUSH_PROMISE frame is sent.283        We may only send PUSH_PROMISE frames if we're a server.284        """285        if self.client is True:286            msg = "Cannot push streams from client peers."287            raise ProtocolError(msg)288 289        event = _PushedRequestSent()290        return [event]291 292    def recv_push_promise(self, previous_state: StreamState) -> list[Event]:293        """294        Fires on the already-existing stream when a PUSH_PROMISE frame is295        received. We may only receive PUSH_PROMISE frames if we're a client.296 297        Fires a PushedStreamReceived event.298        """299        if not self.client:300            if self.client is None:  # pragma: no cover301                msg = "Idle streams cannot receive pushes"302            else:  # pragma: no cover303                msg = "Cannot receive pushed streams as a server"304            raise ProtocolError(msg)305 306        event = PushedStreamReceived()307        event.parent_stream_id = self.stream_id308        return [event]309 310    def send_end_stream(self, previous_state: StreamState) -> list[Event]:311        """312        Called when an attempt is made to send END_STREAM in the313        HALF_CLOSED_REMOTE state.314        """315        self.stream_closed_by = StreamClosedBy.SEND_END_STREAM316        return []317 318    def send_reset_stream(self, previous_state: StreamState) -> list[Event]:319        """320        Called when an attempt is made to send RST_STREAM in a non-closed321        stream state.322        """323        self.stream_closed_by = StreamClosedBy.SEND_RST_STREAM324        return []325 326    def reset_stream_on_error(self, previous_state: StreamState) -> list[Event]:327        """328        Called when we need to forcefully emit another RST_STREAM frame on329        behalf of the state machine.330 331        If this is the first time we've done this, we should also hang an event332        off the StreamClosedError so that the user can be informed. We know333        it's the first time we've done this if the stream is currently in a334        state other than CLOSED.335        """336        self.stream_closed_by = StreamClosedBy.SEND_RST_STREAM337 338        error = StreamClosedError(self.stream_id)339        error._events = [340            StreamReset(341                stream_id=self.stream_id,342                error_code=ErrorCodes.STREAM_CLOSED,343                remote_reset=False,344            ),345        ]346        raise error347 348    def recv_on_closed_stream(self, previous_state: StreamState) -> list[Event]:349        """350        Called when an unexpected frame is received on an already-closed351        stream.352 353        An endpoint that receives an unexpected frame should treat it as354        a stream error or connection error with type STREAM_CLOSED, depending355        on the specific frame. The error handling is done at a higher level:356        this just raises the appropriate error.357        """358        raise StreamClosedError(self.stream_id)359 360    def send_on_closed_stream(self, previous_state: StreamState) -> list[Event]:361        """362        Called when an attempt is made to send data on an already-closed363        stream.364 365        This essentially overrides the standard logic by throwing a366        more-specific error: StreamClosedError. This is a ProtocolError, so it367        matches the standard API of the state machine, but provides more detail368        to the user.369        """370        raise StreamClosedError(self.stream_id)371 372    def recv_push_on_closed_stream(self, previous_state: StreamState) -> list[Event]:373        """374        Called when a PUSH_PROMISE frame is received on a full stop375        stream.376 377        If the stream was closed by us sending a RST_STREAM frame, then we378        presume that the PUSH_PROMISE was in flight when we reset the parent379        stream. Rathen than accept the new stream, we just reset it.380        Otherwise, we should call this a PROTOCOL_ERROR: pushing a stream on a381        naturally closed stream is a real problem because it creates a brand382        new stream that the remote peer now believes exists.383        """384        assert self.stream_closed_by is not None385 386        if self.stream_closed_by == StreamClosedBy.SEND_RST_STREAM:387            raise StreamClosedError(self.stream_id)388        msg = "Attempted to push on closed stream."389        raise ProtocolError(msg)390 391    def send_push_on_closed_stream(self, previous_state: StreamState) -> list[Event]:392        """393        Called when an attempt is made to push on an already-closed stream.394 395        This essentially overrides the standard logic by providing a more396        useful error message. It's necessary because simply indicating that the397        stream is closed is not enough: there is now a new stream that is not398        allowed to be there. The only recourse is to tear the whole connection399        down.400        """401        msg = "Attempted to push on closed stream."402        raise ProtocolError(msg)403 404    def send_informational_response(self, previous_state: StreamState) -> list[Event]:405        """406        Called when an informational header block is sent (that is, a block407        where the :status header has a 1XX value).408 409        Only enforces that these are sent *before* final headers are sent.410        """411        if self.headers_sent:412            msg = "Information response after final response"413            raise ProtocolError(msg)414 415        event = _ResponseSent()416        return [event]417 418    def recv_informational_response(self, previous_state: StreamState) -> list[Event]:419        """420        Called when an informational header block is received (that is, a block421        where the :status header has a 1XX value).422        """423        if self.headers_received:424            msg = "Informational response after final response"425            raise ProtocolError(msg)426 427        return [InformationalResponseReceived(stream_id=self.stream_id)]428 429    def recv_alt_svc(self, previous_state: StreamState) -> list[Event]:430        """431        Called when receiving an ALTSVC frame.432 433        RFC 7838 allows us to receive ALTSVC frames at any stream state, which434        is really absurdly overzealous. For that reason, we want to limit the435        states in which we can actually receive it. It's really only sensible436        to receive it after we've sent our own headers and before the server437        has sent its header block: the server can't guarantee that we have any438        state around after it completes its header block, and the server439        doesn't know what origin we're talking about before we've sent ours.440 441        For that reason, this function applies a few extra checks on both state442        and some of the little state variables we keep around. If those suggest443        an unreasonable situation for the ALTSVC frame to have been sent in,444        we quietly ignore it (as RFC 7838 suggests).445 446        This function is also *not* always called by the state machine. In some447        states (IDLE, RESERVED_LOCAL, CLOSED) we don't bother to call it,448        because we know the frame cannot be valid in that state (IDLE because449        the server cannot know what origin the stream applies to, CLOSED450        because the server cannot assume we still have state around,451        RESERVED_LOCAL because by definition if we're in the RESERVED_LOCAL452        state then *we* are the server).453        """454        # Servers can't receive ALTSVC frames, but RFC 7838 tells us to ignore455        # them.456        if self.client is False:457            return []458 459        # If we've received the response headers from the server they can't460        # guarantee we still have any state around. Other implementations461        # (like nghttp2) ignore ALTSVC in this state, so we will too.462        if self.headers_received:463            return []464 465        # Otherwise, this is a sensible enough frame to have received. Return466        # the event and let it get populated.467        return [AlternativeServiceAvailable()]468 469    def send_alt_svc(self, previous_state: StreamState) -> list[Event]:470        """471        Called when sending an ALTSVC frame on this stream.472 473        For consistency with the restrictions we apply on receiving ALTSVC474        frames in ``recv_alt_svc``, we want to restrict when users can send475        ALTSVC frames to the situations when we ourselves would accept them.476 477        That means: when we are a server, when we have received the request478        headers, and when we have not yet sent our own response headers.479        """480        # We should not send ALTSVC after we've sent response headers, as the481        # client may have disposed of its state.482        if self.headers_sent:483            msg = "Cannot send ALTSVC after sending response headers."484            raise ProtocolError(msg)485        return []486 487 488 489# STATE MACHINE490#491# The stream state machine is defined here to avoid the need to allocate it492# repeatedly for each stream. It cannot be defined in the stream class because493# it needs to be able to reference the callbacks defined on the class, but494# because Python's scoping rules are weird the class object is not actually in495# scope during the body of the class object.496#497# For the sake of clarity, we reproduce the RFC 7540 state machine here:498#499#                          +--------+500#                  send PP |        | recv PP501#                 ,--------|  idle  |--------.502#                /         |        |         \503#               v          +--------+          v504#        +----------+          |           +----------+505#        |          |          | send H /  |          |506# ,------| reserved |          | recv H    | reserved |------.507# |      | (local)  |          |           | (remote) |      |508# |      +----------+          v           +----------+      |509# |          |             +--------+             |          |510# |          |     recv ES |        | send ES     |          |511# |   send H |     ,-------|  open  |-------.     | recv H   |512# |          |    /        |        |        \    |          |513# |          v   v         +--------+         v   v          |514# |      +----------+          |           +----------+      |515# |      |   half   |          |           |   half   |      |516# |      |  closed  |          | send R /  |  closed  |      |517# |      | (remote) |          | recv R    | (local)  |      |518# |      +----------+          |           +----------+      |519# |           |                |                 |           |520# |           | send ES /      |       recv ES / |           |521# |           | send R /       v        send R / |           |522# |           | recv R     +--------+   recv R   |           |523# | send R /  `----------->|        |<-----------'  send R / |524# | recv R                 | closed |               recv R   |525# `----------------------->|        |<----------------------'526#                          +--------+527#528#    send:   endpoint sends this frame529#    recv:   endpoint receives this frame530#531#    H:  HEADERS frame (with implied CONTINUATIONs)532#    PP: PUSH_PROMISE frame (with implied CONTINUATIONs)533#    ES: END_STREAM flag534#    R:  RST_STREAM frame535#536# For the purposes of this state machine we treat HEADERS and their537# associated CONTINUATION frames as a single jumbo frame. The protocol538# allows/requires this by preventing other frames from being interleved in539# between HEADERS/CONTINUATION frames. However, if a CONTINUATION frame is540# received without a prior HEADERS frame, it *will* be passed to this state541# machine. The state machine should always reject that frame, either as an542# invalid transition or because the stream is closed.543#544# There is a confusing relationship around PUSH_PROMISE frames. The state545# machine above considers them to be frames belonging to the new stream,546# which is *somewhat* true. However, they are sent with the stream ID of547# their related stream, and are only sendable in some cases.548# For this reason, our state machine implementation below allows for549# PUSH_PROMISE frames both in the IDLE state (as in the diagram), but also550# in the OPEN, HALF_CLOSED_LOCAL, and HALF_CLOSED_REMOTE states.551# Essentially, for h2, PUSH_PROMISE frames are effectively sent on552# two streams.553#554# The _transitions dictionary contains a mapping of tuples of555# (state, input) to tuples of (side_effect_function, end_state). This556# map contains all allowed transitions: anything not in this map is557# invalid and immediately causes a transition to ``closed``.558_transitions: dict[559    tuple[StreamState, StreamInputs],560    tuple[Callable[[H2StreamStateMachine, StreamState], list[Event]] | None, StreamState],561] = {562    # State: idle563    (StreamState.IDLE, StreamInputs.SEND_HEADERS):564        (H2StreamStateMachine.request_sent, StreamState.OPEN),565    (StreamState.IDLE, StreamInputs.RECV_HEADERS):566        (H2StreamStateMachine.request_received, StreamState.OPEN),567    (StreamState.IDLE, StreamInputs.RECV_DATA):568        (H2StreamStateMachine.reset_stream_on_error, StreamState.CLOSED),569    (StreamState.IDLE, StreamInputs.SEND_PUSH_PROMISE):570        (H2StreamStateMachine.send_new_pushed_stream,571            StreamState.RESERVED_LOCAL),572    (StreamState.IDLE, StreamInputs.RECV_PUSH_PROMISE):573        (H2StreamStateMachine.recv_new_pushed_stream,574            StreamState.RESERVED_REMOTE),575    (StreamState.IDLE, StreamInputs.RECV_ALTERNATIVE_SERVICE):576        (None, StreamState.IDLE),577    (StreamState.IDLE, StreamInputs.UPGRADE_CLIENT):578        (H2StreamStateMachine.request_sent, StreamState.HALF_CLOSED_LOCAL),579    (StreamState.IDLE, StreamInputs.UPGRADE_SERVER):580        (H2StreamStateMachine.request_received,581            StreamState.HALF_CLOSED_REMOTE),582 583    # State: reserved local584    (StreamState.RESERVED_LOCAL, StreamInputs.SEND_HEADERS):585        (H2StreamStateMachine.response_sent, StreamState.HALF_CLOSED_REMOTE),586    (StreamState.RESERVED_LOCAL, StreamInputs.RECV_DATA):587        (H2StreamStateMachine.reset_stream_on_error, StreamState.CLOSED),588    (StreamState.RESERVED_LOCAL, StreamInputs.SEND_WINDOW_UPDATE):589        (None, StreamState.RESERVED_LOCAL),590    (StreamState.RESERVED_LOCAL, StreamInputs.RECV_WINDOW_UPDATE):591        (H2StreamStateMachine.window_updated, StreamState.RESERVED_LOCAL),592    (StreamState.RESERVED_LOCAL, StreamInputs.SEND_RST_STREAM):593        (H2StreamStateMachine.send_reset_stream, StreamState.CLOSED),594    (StreamState.RESERVED_LOCAL, StreamInputs.RECV_RST_STREAM):595        (H2StreamStateMachine.stream_reset, StreamState.CLOSED),596    (StreamState.RESERVED_LOCAL, StreamInputs.SEND_ALTERNATIVE_SERVICE):597        (H2StreamStateMachine.send_alt_svc, StreamState.RESERVED_LOCAL),598    (StreamState.RESERVED_LOCAL, StreamInputs.RECV_ALTERNATIVE_SERVICE):599        (None, StreamState.RESERVED_LOCAL),600 601    # State: reserved remote602    (StreamState.RESERVED_REMOTE, StreamInputs.RECV_HEADERS):603        (H2StreamStateMachine.response_received,604            StreamState.HALF_CLOSED_LOCAL),605    (StreamState.RESERVED_REMOTE, StreamInputs.RECV_DATA):606        (H2StreamStateMachine.reset_stream_on_error, StreamState.CLOSED),607    (StreamState.RESERVED_REMOTE, StreamInputs.SEND_WINDOW_UPDATE):608        (None, StreamState.RESERVED_REMOTE),609    (StreamState.RESERVED_REMOTE, StreamInputs.RECV_WINDOW_UPDATE):610        (H2StreamStateMachine.window_updated, StreamState.RESERVED_REMOTE),611    (StreamState.RESERVED_REMOTE, StreamInputs.SEND_RST_STREAM):612        (H2StreamStateMachine.send_reset_stream, StreamState.CLOSED),613    (StreamState.RESERVED_REMOTE, StreamInputs.RECV_RST_STREAM):614        (H2StreamStateMachine.stream_reset, StreamState.CLOSED),615    (StreamState.RESERVED_REMOTE, StreamInputs.RECV_ALTERNATIVE_SERVICE):616        (H2StreamStateMachine.recv_alt_svc, StreamState.RESERVED_REMOTE),617 618    # State: open619    (StreamState.OPEN, StreamInputs.SEND_HEADERS):620        (H2StreamStateMachine.response_sent, StreamState.OPEN),621    (StreamState.OPEN, StreamInputs.RECV_HEADERS):622        (H2StreamStateMachine.response_received, StreamState.OPEN),623    (StreamState.OPEN, StreamInputs.SEND_DATA):624        (None, StreamState.OPEN),625    (StreamState.OPEN, StreamInputs.RECV_DATA):626        (H2StreamStateMachine.data_received, StreamState.OPEN),627    (StreamState.OPEN, StreamInputs.SEND_END_STREAM):628        (None, StreamState.HALF_CLOSED_LOCAL),629    (StreamState.OPEN, StreamInputs.RECV_END_STREAM):630        (H2StreamStateMachine.stream_half_closed,631         StreamState.HALF_CLOSED_REMOTE),632    (StreamState.OPEN, StreamInputs.SEND_WINDOW_UPDATE):633        (None, StreamState.OPEN),634    (StreamState.OPEN, StreamInputs.RECV_WINDOW_UPDATE):635        (H2StreamStateMachine.window_updated, StreamState.OPEN),636    (StreamState.OPEN, StreamInputs.SEND_RST_STREAM):637        (H2StreamStateMachine.send_reset_stream, StreamState.CLOSED),638    (StreamState.OPEN, StreamInputs.RECV_RST_STREAM):639        (H2StreamStateMachine.stream_reset, StreamState.CLOSED),640    (StreamState.OPEN, StreamInputs.SEND_PUSH_PROMISE):641        (H2StreamStateMachine.send_push_promise, StreamState.OPEN),642    (StreamState.OPEN, StreamInputs.RECV_PUSH_PROMISE):643        (H2StreamStateMachine.recv_push_promise, StreamState.OPEN),644    (StreamState.OPEN, StreamInputs.SEND_INFORMATIONAL_HEADERS):645        (H2StreamStateMachine.send_informational_response, StreamState.OPEN),646    (StreamState.OPEN, StreamInputs.RECV_INFORMATIONAL_HEADERS):647        (H2StreamStateMachine.recv_informational_response, StreamState.OPEN),648    (StreamState.OPEN, StreamInputs.SEND_ALTERNATIVE_SERVICE):649        (H2StreamStateMachine.send_alt_svc, StreamState.OPEN),650    (StreamState.OPEN, StreamInputs.RECV_ALTERNATIVE_SERVICE):651        (H2StreamStateMachine.recv_alt_svc, StreamState.OPEN),652 653    # State: half-closed remote654    (StreamState.HALF_CLOSED_REMOTE, StreamInputs.SEND_HEADERS):655        (H2StreamStateMachine.response_sent, StreamState.HALF_CLOSED_REMOTE),656    (StreamState.HALF_CLOSED_REMOTE, StreamInputs.RECV_HEADERS):657        (H2StreamStateMachine.reset_stream_on_error, StreamState.CLOSED),658    (StreamState.HALF_CLOSED_REMOTE, StreamInputs.SEND_DATA):659        (None, StreamState.HALF_CLOSED_REMOTE),660    (StreamState.HALF_CLOSED_REMOTE, StreamInputs.RECV_DATA):661        (H2StreamStateMachine.reset_stream_on_error, StreamState.CLOSED),662    (StreamState.HALF_CLOSED_REMOTE, StreamInputs.SEND_END_STREAM):663        (H2StreamStateMachine.send_end_stream, StreamState.CLOSED),664    (StreamState.HALF_CLOSED_REMOTE, StreamInputs.SEND_WINDOW_UPDATE):665        (None, StreamState.HALF_CLOSED_REMOTE),666    (StreamState.HALF_CLOSED_REMOTE, StreamInputs.RECV_WINDOW_UPDATE):667        (H2StreamStateMachine.window_updated, StreamState.HALF_CLOSED_REMOTE),668    (StreamState.HALF_CLOSED_REMOTE, StreamInputs.SEND_RST_STREAM):669        (H2StreamStateMachine.send_reset_stream, StreamState.CLOSED),670    (StreamState.HALF_CLOSED_REMOTE, StreamInputs.RECV_RST_STREAM):671        (H2StreamStateMachine.stream_reset, StreamState.CLOSED),672    (StreamState.HALF_CLOSED_REMOTE, StreamInputs.SEND_PUSH_PROMISE):673        (H2StreamStateMachine.send_push_promise,674            StreamState.HALF_CLOSED_REMOTE),675    (StreamState.HALF_CLOSED_REMOTE, StreamInputs.RECV_PUSH_PROMISE):676        (H2StreamStateMachine.reset_stream_on_error, StreamState.CLOSED),677    (StreamState.HALF_CLOSED_REMOTE, StreamInputs.SEND_INFORMATIONAL_HEADERS):678        (H2StreamStateMachine.send_informational_response,679            StreamState.HALF_CLOSED_REMOTE),680    (StreamState.HALF_CLOSED_REMOTE, StreamInputs.SEND_ALTERNATIVE_SERVICE):681        (H2StreamStateMachine.send_alt_svc, StreamState.HALF_CLOSED_REMOTE),682    (StreamState.HALF_CLOSED_REMOTE, StreamInputs.RECV_ALTERNATIVE_SERVICE):683        (H2StreamStateMachine.recv_alt_svc, StreamState.HALF_CLOSED_REMOTE),684 685    # State: half-closed local686    (StreamState.HALF_CLOSED_LOCAL, StreamInputs.RECV_HEADERS):687        (H2StreamStateMachine.response_received,688            StreamState.HALF_CLOSED_LOCAL),689    (StreamState.HALF_CLOSED_LOCAL, StreamInputs.RECV_DATA):690        (H2StreamStateMachine.data_received, StreamState.HALF_CLOSED_LOCAL),691    (StreamState.HALF_CLOSED_LOCAL, StreamInputs.RECV_END_STREAM):692        (H2StreamStateMachine.stream_ended, StreamState.CLOSED),693    (StreamState.HALF_CLOSED_LOCAL, StreamInputs.SEND_WINDOW_UPDATE):694        (None, StreamState.HALF_CLOSED_LOCAL),695    (StreamState.HALF_CLOSED_LOCAL, StreamInputs.RECV_WINDOW_UPDATE):696        (H2StreamStateMachine.window_updated, StreamState.HALF_CLOSED_LOCAL),697    (StreamState.HALF_CLOSED_LOCAL, StreamInputs.SEND_RST_STREAM):698        (H2StreamStateMachine.send_reset_stream, StreamState.CLOSED),699    (StreamState.HALF_CLOSED_LOCAL, StreamInputs.RECV_RST_STREAM):700        (H2StreamStateMachine.stream_reset, StreamState.CLOSED),701    (StreamState.HALF_CLOSED_LOCAL, StreamInputs.RECV_PUSH_PROMISE):702        (H2StreamStateMachine.recv_push_promise,703            StreamState.HALF_CLOSED_LOCAL),704    (StreamState.HALF_CLOSED_LOCAL, StreamInputs.RECV_INFORMATIONAL_HEADERS):705        (H2StreamStateMachine.recv_informational_response,706            StreamState.HALF_CLOSED_LOCAL),707    (StreamState.HALF_CLOSED_LOCAL, StreamInputs.SEND_ALTERNATIVE_SERVICE):708        (H2StreamStateMachine.send_alt_svc, StreamState.HALF_CLOSED_LOCAL),709    (StreamState.HALF_CLOSED_LOCAL, StreamInputs.RECV_ALTERNATIVE_SERVICE):710        (H2StreamStateMachine.recv_alt_svc, StreamState.HALF_CLOSED_LOCAL),711 712    # State: closed713    (StreamState.CLOSED, StreamInputs.RECV_END_STREAM):714        (None, StreamState.CLOSED),715    (StreamState.CLOSED, StreamInputs.RECV_ALTERNATIVE_SERVICE):716        (None, StreamState.CLOSED),717 718    # RFC 7540 Section 5.1 defines how the end point should react when719    # receiving a frame on a closed stream with the following statements:720    #721    # > An endpoint that receives any frame other than PRIORITY after receiving722    # > a RST_STREAM MUST treat that as a stream error of type STREAM_CLOSED.723    # > An endpoint that receives any frames after receiving a frame with the724    # > END_STREAM flag set MUST treat that as a connection error of type725    # > STREAM_CLOSED.726    (StreamState.CLOSED, StreamInputs.RECV_HEADERS):727        (H2StreamStateMachine.recv_on_closed_stream, StreamState.CLOSED),728    (StreamState.CLOSED, StreamInputs.RECV_DATA):729        (H2StreamStateMachine.recv_on_closed_stream, StreamState.CLOSED),730 731    # > WINDOW_UPDATE or RST_STREAM frames can be received in this state732    # > for a short period after a DATA or HEADERS frame containing a733    # > END_STREAM flag is sent, as instructed in RFC 7540 Section 5.1. But we734    # > don't have access to a clock so we just always allow it.735    (StreamState.CLOSED, StreamInputs.RECV_WINDOW_UPDATE):736        (None, StreamState.CLOSED),737    (StreamState.CLOSED, StreamInputs.RECV_RST_STREAM):738        (None, StreamState.CLOSED),739 740    # > A receiver MUST treat the receipt of a PUSH_PROMISE on a stream that is741    # > neither "open" nor "half-closed (local)" as a connection error of type742    # > PROTOCOL_ERROR.743    (StreamState.CLOSED, StreamInputs.RECV_PUSH_PROMISE):744        (H2StreamStateMachine.recv_push_on_closed_stream, StreamState.CLOSED),745 746    # Also, users should be forbidden from sending on closed streams.747    (StreamState.CLOSED, StreamInputs.SEND_HEADERS):748        (H2StreamStateMachine.send_on_closed_stream, StreamState.CLOSED),749    (StreamState.CLOSED, StreamInputs.SEND_PUSH_PROMISE):750        (H2StreamStateMachine.send_push_on_closed_stream, StreamState.CLOSED),751    (StreamState.CLOSED, StreamInputs.SEND_RST_STREAM):752        (H2StreamStateMachine.send_on_closed_stream, StreamState.CLOSED),753    (StreamState.CLOSED, StreamInputs.SEND_DATA):754        (H2StreamStateMachine.send_on_closed_stream, StreamState.CLOSED),755    (StreamState.CLOSED, StreamInputs.SEND_WINDOW_UPDATE):756        (H2StreamStateMachine.send_on_closed_stream, StreamState.CLOSED),757    (StreamState.CLOSED, StreamInputs.SEND_END_STREAM):758        (H2StreamStateMachine.send_on_closed_stream, StreamState.CLOSED),759}760 761 762class H2Stream:763    """764    A low-level HTTP/2 stream object. This handles building and receiving765    frames and maintains per-stream state.766 767    This wraps a HTTP/2 Stream state machine implementation, ensuring that768    frames can only be sent/received when the stream is in a valid state.769    Attempts to create frames that cannot be sent will raise a770    ``ProtocolError``.771    """772 773    def __init__(self,774                 stream_id: int,775                 config: H2Configuration,776                 inbound_window_size: int,777                 outbound_window_size: int) -> None:778        self.state_machine = H2StreamStateMachine(stream_id)779        self.stream_id = stream_id780        self.max_outbound_frame_size: int | None = None781        self.request_method: bytes | None = None782 783        # The current value of the outbound stream flow control window784        self.outbound_flow_control_window = outbound_window_size785 786        # The flow control manager.787        self._inbound_window_manager = WindowManager(inbound_window_size)788 789        # The expected content length, if any.790        self._expected_content_length: int | None = None791 792        # The actual received content length. Always tracked.793        self._actual_content_length = 0794 795        # The authority we believe this stream belongs to.796        self._authority: bytes | None = None797 798        # The configuration for this stream.799        self.config = config800 801    def __repr__(self) -> str:802        return f"<{type(self).__name__} id:{self.stream_id} state:{self.state_machine.state!r}>"803 804    @property805    def inbound_flow_control_window(self) -> int:806        """807        The size of the inbound flow control window for the stream. This is808        rarely publicly useful: instead, use :meth:`remote_flow_control_window809        <h2.stream.H2Stream.remote_flow_control_window>`. This shortcut is810        largely present to provide a shortcut to this data.811        """812        return self._inbound_window_manager.current_window_size813 814    @property815    def open(self) -> bool:816        """817        Whether the stream is 'open' in any sense: that is, whether it counts818        against the number of concurrent streams.819        """820        # RFC 7540 Section 5.1.2 defines 'open' for this purpose to mean either821        # the OPEN state or either of the HALF_CLOSED states. Perplexingly,822        # this excludes the reserved states.823        # For more detail on why we're doing this in this slightly weird way,824        # see the comment on ``STREAM_OPEN`` at the top of the file.825        return STREAM_OPEN[self.state_machine.state]826 827    @property828    def closed(self) -> bool:829        """830        Whether the stream is closed.831        """832        return self.state_machine.state == StreamState.CLOSED833 834    @property835    def closed_by(self) -> StreamClosedBy | None:836        """837        Returns how the stream was closed, as one of StreamClosedBy.838        """839        return self.state_machine.stream_closed_by840 841    def upgrade(self, client_side: bool) -> None:842        """843        Called by the connection to indicate that this stream is the initial844        request/response of an upgraded connection. Places the stream into an845        appropriate state.846        """847        self.config.logger.debug("Upgrading %r", self)848 849        assert self.stream_id == 1850        input_ = (851            StreamInputs.UPGRADE_CLIENT if client_side852            else StreamInputs.UPGRADE_SERVER853        )854 855        # This may return events, we deliberately don't want them.856        self.state_machine.process_input(input_)857 858    def send_headers(self,859                     headers: Iterable[HeaderWeaklyTyped],860                     encoder: Encoder,861                     end_stream: bool = False) -> list[HeadersFrame | ContinuationFrame | PushPromiseFrame]:862        """863        Returns a list of HEADERS/CONTINUATION frames to emit as either headers864        or trailers.865        """866        self.config.logger.debug("Send headers %s on %r", headers, self)867 868        # Because encoding headers makes an irreversible change to the header869        # compression context, we make the state transition before we encode870        # them.871 872        # First, check if we're a client. If we are, no problem: if we aren't,873        # we need to scan the header block to see if this is an informational874        # response.875        input_ = StreamInputs.SEND_HEADERS876 877        bytes_headers = utf8_encode_headers(headers)878 879        if ((not self.state_machine.client) and880                is_informational_response(bytes_headers)):881            if end_stream:882                msg = "Cannot set END_STREAM on informational responses."883                raise ProtocolError(msg)884 885            input_ = StreamInputs.SEND_INFORMATIONAL_HEADERS886 887        events = self.state_machine.process_input(input_)888 889        hf = HeadersFrame(self.stream_id)890        hdr_validation_flags = self._build_hdr_validation_flags(events)891        frames = self._build_headers_frames(892            bytes_headers, encoder, hf, hdr_validation_flags,893        )894 895        if end_stream:896            # Not a bug: the END_STREAM flag is valid on the initial HEADERS897            # frame, not the CONTINUATION frames that follow.898            self.state_machine.process_input(StreamInputs.SEND_END_STREAM)899            frames[0].flags.add("END_STREAM")900 901        if self.state_machine.trailers_sent and not end_stream:902            msg = "Trailers must have END_STREAM set."903            raise ProtocolError(msg)904 905        if self.state_machine.client and self._authority is None:906            self._authority = authority_from_headers(bytes_headers)907 908        # store request method for _initialize_content_length909        self.request_method = extract_method_header(bytes_headers)910 911        return frames912 913    def push_stream_in_band(self,914                            related_stream_id: int,915                            headers: Iterable[HeaderWeaklyTyped],916                            encoder: Encoder) -> list[HeadersFrame | ContinuationFrame | PushPromiseFrame]:917        """918        Returns a list of PUSH_PROMISE/CONTINUATION frames to emit as a pushed919        stream header. Called on the stream that has the PUSH_PROMISE frame920        sent on it.921        """922        self.config.logger.debug("Push stream %r", self)923 924        # Because encoding headers makes an irreversible change to the header925        # compression context, we make the state transition *first*.926 927        events = self.state_machine.process_input(928            StreamInputs.SEND_PUSH_PROMISE,929        )930 931        ppf = PushPromiseFrame(self.stream_id)932        ppf.promised_stream_id = related_stream_id933        hdr_validation_flags = self._build_hdr_validation_flags(events)934 935        bytes_headers = utf8_encode_headers(headers)936 937        return self._build_headers_frames(938            bytes_headers, encoder, ppf, hdr_validation_flags,939        )940 941 942    def locally_pushed(self) -> list[Frame]:943        """944        Mark this stream as one that was pushed by this peer. Must be called945        immediately after initialization. Sends no frames, simply updates the946        state machine.947        """948        # This does not trigger any events.949        events = self.state_machine.process_input(950            StreamInputs.SEND_PUSH_PROMISE,951        )952        assert not events953        return []954 955    def send_data(self,956                  data: bytes | memoryview,957                  end_stream: bool = False,958                  pad_length: int | None = None) -> list[Frame]:959        """960        Prepare some data frames. Optionally end the stream.961 962        .. warning:: Does not perform flow control checks.963        """964        self.config.logger.debug(965            "Send data on %r with end stream set to %s", self, end_stream,966        )967 968        self.state_machine.process_input(StreamInputs.SEND_DATA)969 970        df = DataFrame(self.stream_id)971        df.data = data972        if end_stream:973            self.state_machine.process_input(StreamInputs.SEND_END_STREAM)974            df.flags.add("END_STREAM")975        if pad_length is not None:976            df.flags.add("PADDED")977            df.pad_length = pad_length978 979        # Subtract flow_controlled_length to account for possible padding980        self.outbound_flow_control_window -= df.flow_controlled_length981        assert self.outbound_flow_control_window >= 0982 983        return [df]984 985    def end_stream(self) -> list[Frame]:986        """987        End a stream without sending data.988        """989        self.config.logger.debug("End stream %r", self)990 991        self.state_machine.process_input(StreamInputs.SEND_END_STREAM)992        df = DataFrame(self.stream_id)993        df.flags.add("END_STREAM")994        return [df]995 996    def advertise_alternative_service(self, field_value: bytes) -> list[Frame]:997        """998        Advertise an RFC 7838 alternative service. The semantics of this are999        better documented in the ``H2Connection`` class.1000        """1001        self.config.logger.debug(1002            "Advertise alternative service of %r for %r", field_value, self,1003        )1004        self.state_machine.process_input(StreamInputs.SEND_ALTERNATIVE_SERVICE)1005        asf = AltSvcFrame(self.stream_id)1006        asf.field = field_value1007        return [asf]1008 1009    def increase_flow_control_window(self, increment: int) -> list[Frame]:1010        """1011        Increase the size of the flow control window for the remote side.1012        """1013        self.config.logger.debug(1014            "Increase flow control window for %r by %d",1015            self, increment,1016        )1017        self.state_machine.process_input(StreamInputs.SEND_WINDOW_UPDATE)1018        self._inbound_window_manager.window_opened(increment)1019 1020        wuf = WindowUpdateFrame(self.stream_id)1021        wuf.window_increment = increment1022        return [wuf]1023 1024    def receive_push_promise_in_band(self,1025                                     promised_stream_id: int,1026                                     headers: Iterable[Header],1027                                     header_encoding: bool | str | None) -> tuple[list[Frame], list[Event]]:1028        """1029        Receives a push promise frame sent on this stream, pushing a remote1030        stream. This is called on the stream that has the PUSH_PROMISE sent1031        on it.1032        """1033        self.config.logger.debug(1034            "Receive Push Promise on %r for remote stream %d",1035            self, promised_stream_id,1036        )1037        events = self.state_machine.process_input(1038            StreamInputs.RECV_PUSH_PROMISE,1039        )1040        push_event = cast("PushedStreamReceived", events[0])1041        push_event.pushed_stream_id = promised_stream_id1042 1043        hdr_validation_flags = self._build_hdr_validation_flags(events)1044        push_event.headers = self._process_received_headers(1045            headers, hdr_validation_flags, header_encoding,1046        )1047        return [], events1048 1049    def remotely_pushed(self, pushed_headers: Iterable[Header]) -> tuple[list[Frame], list[Event]]:1050        """1051        Mark this stream as one that was pushed by the remote peer. Must be1052        called immediately after initialization. Sends no frames, simply1053        updates the state machine.1054        """1055        self.config.logger.debug("%r pushed by remote peer", self)1056        events = self.state_machine.process_input(1057            StreamInputs.RECV_PUSH_PROMISE,1058        )1059        self._authority = authority_from_headers(pushed_headers)1060        return [], events1061 1062    def receive_headers(self,1063                        headers: Iterable[Header],1064                        end_stream: bool,1065                        header_encoding: bool | str | None) -> tuple[list[Frame], list[Event]]:1066        """1067        Receive a set of headers (or trailers).1068        """1069        if is_informational_response(headers):1070            if end_stream:1071                msg = "Cannot set END_STREAM on informational responses"1072                raise ProtocolError(msg)1073            input_ = StreamInputs.RECV_INFORMATIONAL_HEADERS1074        else:1075            input_ = StreamInputs.RECV_HEADERS1076 1077        events = self.state_machine.process_input(input_)1078        headers_event = cast(1079            "Union[RequestReceived, ResponseReceived, TrailersReceived, InformationalResponseReceived]",1080            events[0],1081        )1082 1083        if end_stream:1084            es_events = self.state_machine.process_input(1085                StreamInputs.RECV_END_STREAM,1086            )1087            # We ensured it's not an information response at the beginning of the method.1088            cast(1089                "Union[RequestReceived, ResponseReceived, TrailersReceived]",1090                headers_event,1091            ).stream_ended = cast("StreamEnded", es_events[0])1092            events += es_events1093 1094        self._initialize_content_length(headers)1095 1096        if isinstance(headers_event, TrailersReceived) and not end_stream:1097            msg = "Trailers must have END_STREAM set"1098            raise ProtocolError(msg)1099 1100        hdr_validation_flags = self._build_hdr_validation_flags(events)1101        headers_event.headers = self._process_received_headers(1102            headers, hdr_validation_flags, header_encoding,1103        )1104        return [], events1105 1106    def receive_data(self, data: bytes, end_stream: bool, flow_control_len: int) -> tuple[list[Frame], list[Event]]:1107        """1108        Receive some data.1109        """1110        self.config.logger.debug(1111            "Receive data on %r with end stream %s and flow control length "1112            "set to %d", self, end_stream, flow_control_len,1113        )1114        events = self.state_machine.process_input(StreamInputs.RECV_DATA)1115        data_event = cast("DataReceived", events[0])1116        self._inbound_window_manager.window_consumed(flow_control_len)1117        self._track_content_length(len(data), end_stream)1118 1119        if end_stream:1120            es_events = self.state_machine.process_input(1121                StreamInputs.RECV_END_STREAM,1122            )1123            data_event.stream_ended = cast("StreamEnded", es_events[0])1124            events.extend(es_events)1125 1126        data_event.data = data1127        data_event.flow_controlled_length = flow_control_len1128        return [], events1129 1130    def receive_window_update(self, increment: int) -> tuple[list[Frame], list[Event]]:1131        """1132        Handle a WINDOW_UPDATE increment.1133        """1134        self.config.logger.debug(1135            "Receive Window Update on %r for increment of %d",1136            self, increment,1137        )1138        events = self.state_machine.process_input(1139            StreamInputs.RECV_WINDOW_UPDATE,1140        )1141        frames = []1142 1143        # If we encounter a problem with incrementing the flow control window,1144        # this should be treated as a *stream* error, not a *connection* error.1145        # That means we need to catch the error and forcibly close the stream.1146        if events:1147            cast("WindowUpdated", events[0]).delta = increment1148            try:1149                self.outbound_flow_control_window = guard_increment_window(1150                    self.outbound_flow_control_window,1151                    increment,1152                )1153            except FlowControlError:1154                # Ok, this is bad. We're going to need to perform a local1155                # reset.1156                events = [1157                    StreamReset(1158                        stream_id=self.stream_id,1159                        error_code=ErrorCodes.FLOW_CONTROL_ERROR,1160                        remote_reset=False,1161                    ),1162                ]1163                frames = self.reset_stream(ErrorCodes.FLOW_CONTROL_ERROR)1164 1165        return frames, events1166 1167    def receive_continuation(self) -> None:1168        """1169        A naked CONTINUATION frame has been received. This is always an error,1170        but the type of error it is depends on the state of the stream and must1171        transition the state of the stream, so we need to handle it.1172        """1173        self.config.logger.debug("Receive Continuation frame on %r", self)1174        self.state_machine.process_input(1175            StreamInputs.RECV_CONTINUATION,1176        )1177        msg = "Should not be reachable"  # pragma: no cover1178        raise AssertionError(msg)  # pragma: no cover1179 1180    def receive_alt_svc(self, frame: AltSvcFrame) -> tuple[list[Frame], list[Event]]:1181        """1182        An Alternative Service frame was received on the stream. This frame1183        inherits the origin associated with this stream.1184        """1185        self.config.logger.debug(1186            "Receive Alternative Service frame on stream %r", self,1187        )1188 1189        # If the origin is present, RFC 7838 says we have to ignore it.1190        if frame.origin:1191            return [], []1192 1193        events = self.state_machine.process_input(1194            StreamInputs.RECV_ALTERNATIVE_SERVICE,1195        )1196 1197        # There are lots of situations where we want to ignore the ALTSVC1198        # frame. If we need to pay attention, we'll have an event and should1199        # fill it out.1200        if events:

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

codekingpro/portable-devtools · Team Ai