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