diff --git a/scrapy/core/http2/protocol.py b/scrapy/core/http2/protocol.py index d6134183a..188c14c15 100644 --- a/scrapy/core/http2/protocol.py +++ b/scrapy/core/http2/protocol.py @@ -11,7 +11,7 @@ from h2.events import ( ) from twisted.internet.protocol import connectionDone, Protocol -from scrapy.core.http2.stream import Stream +from scrapy.core.http2.stream import Stream, StreamCloseReason from scrapy.http import Request LOGGER = logging.getLogger(__name__) @@ -71,7 +71,7 @@ class H2ClientProtocol(Protocol): stream_id=stream_id, request=request, connection=self.conn, - metadata=self._metadata, + conn_metadata=self._metadata, write_to_transport=self._write_to_transport, cb_close=self._stream_close_cb ) @@ -133,7 +133,7 @@ class H2ClientProtocol(Protocol): """ # Pop all streams which were pending and were not yet started for stream_id in list(self.streams): - self.streams[stream_id].close() + self.streams[stream_id].close(StreamCloseReason.CONNECTION_LOST) self.conn.close_connection() @@ -180,19 +180,18 @@ class H2ClientProtocol(Protocol): def stream_ended(self, event: StreamEnded): stream_id = event.stream_id - self.streams[stream_id].close() + self.streams[stream_id].close(StreamCloseReason.ENDED) def stream_reset(self, event: StreamReset): # TODO: event.stream_id was abruptly closed # Q. What should be the response? (Failure/Partial/???) - self.streams[event.stream_id].close(event) + self.streams[event.stream_id].close(StreamCloseReason.RESET) def window_updated(self, event: WindowUpdated): stream_id = event.stream_id if stream_id != 0: - self.streams[stream_id].receive_window_update(event.delta) + self.streams[stream_id].receive_window_update() else: # Send leftover data for all the streams for stream in self.streams.values(): - if stream.request_sent: - stream.send_data() + stream.receive_window_update() diff --git a/scrapy/core/http2/stream.py b/scrapy/core/http2/stream.py index e8b4471d6..07a4428c8 100644 --- a/scrapy/core/http2/stream.py +++ b/scrapy/core/http2/stream.py @@ -1,11 +1,15 @@ import logging +from enum import IntFlag, auto +from io import BytesIO from typing import Dict from urllib.parse import urlparse from h2.connection import H2Connection -from h2.events import StreamEnded +from h2.errors import ErrorCodes from h2.exceptions import StreamClosedError -from twisted.internet.defer import Deferred +from twisted.internet.defer import Deferred, CancelledError +from twisted.python.failure import Failure +from twisted.web.client import ResponseFailed from scrapy.http import Request from scrapy.http.headers import Headers @@ -14,6 +18,20 @@ from scrapy.responsetypes import responsetypes LOGGER = logging.getLogger(__name__) +class StreamCloseReason(IntFlag): + # Received a StreamEnded event + ENDED = auto() + + # Received a StreamReset event -- ended abruptly + RESET = auto() + + # Transport connection was lost + CONNECTION_LOST = auto() + + # Expected response body size is more than allowed limit + MAXSIZE_EXCEEDED = auto() + + class Stream: """Represents a single HTTP/2 Stream. @@ -26,13 +44,16 @@ class Stream: """ def __init__( - self, - stream_id: int, - request: Request, - connection: H2Connection, - metadata: Dict, - write_to_transport, - cb_close + self, + stream_id: int, + request: Request, + connection: H2Connection, + conn_metadata: Dict, + write_to_transport, + cb_close, + download_maxsize=0, + download_warnsize=0, + fail_on_data_loss=True ): """ Arguments: @@ -40,7 +61,7 @@ class Stream: uniquely identified by a single integer request {Request} -- HTTP request connection {H2Connection} -- HTTP/2 connection this stream belongs to. - metadata {Dict} -- Reference to dictionary having metadata of HTTP/2 connection + conn_metadata {Dict} -- Reference to dictionary having metadata of HTTP/2 connection write_to_transport {callable} -- Method used to write & send data to the server This method should be used whenever some frame is to be sent to the server. cb_close {callable} -- Method called when this stream is closed @@ -49,16 +70,24 @@ class Stream: self.stream_id = stream_id self._request = request self._conn = connection - self._metadata = metadata + self._conn_metadata = conn_metadata self._write_to_transport = write_to_transport self._cb_close = cb_close - self._request_body = self._request.body - self.content_length = 0 if self._request_body is None else len(self._request_body) + self._download_maxsize = self._request.meta.get('download_maxsize', download_maxsize) + self._download_warnsize = self._request.meta.get('download_warnsize', download_warnsize) + self._fail_on_dataloss = self._request.meta.get('download_fail_on_dataloss', fail_on_data_loss) + + self.request_start_time = None + + self.content_length = 0 if self._request.body is None else len(self._request.body) # Flag to keep track whether this stream has initiated the request self.request_sent = False + # Flag to track whether we have logged about exceeding download warnsize + self._reached_warnsize = False + # Each time we send a data frame, we will decrease value by the amount send. self.remaining_content_length = self.content_length @@ -68,17 +97,17 @@ class Stream: # Flag to keep track whether the server has closed the stream self.stream_closed_server = False - # The amount of data received that counts against the flow control - # window - self._response_flow_controlled_size = 0 - # Private variable used to build the response # this response is then converted to appropriate Response class # passed to the response deferred callback self._response = { # Data received frame by frame from the server is appended # and passed to the response Deferred when completely received. - 'body': b'', + 'body': BytesIO(), + + # The amount of data received that counts against the flow control + # window + 'flow_controlled_size': 0, # Headers received after sending the request 'headers': Headers({}) @@ -137,16 +166,8 @@ class Stream: Warning: Only call this method when stream not closed from client side and has initiated request already by sending HEADER frame. If not then - stream will be closed from client side with 499 response. - - TODO: Q. Should we instead raise ProtocolError here with a proper message? + stream will raise ProtocolError (raise by h2 state machine). """ - if self.stream_closed_local or self.stream_closed_server: - raise StreamClosedError(self.stream_id) - elif not self.request_sent: - self.close() - return - # TODO: # 1. Add test for sending very large data # 2. Add test for small data @@ -154,6 +175,9 @@ class Stream: # 3.1 Large number of request # 3.2 Small number of requests + if self.stream_closed_local: + raise StreamClosedError(self.stream_id) + # Firstly, check what the flow control window is for current stream. window_size = self._conn.local_flow_control_window(stream_id=self.stream_id) @@ -170,7 +194,7 @@ class Stream: chunk_size = min(bytes_to_send_size, max_frame_size) data_chunk_start_id = self.content_length - self.remaining_content_length - data_chunk = self._request_body[data_chunk_start_id:data_chunk_start_id + chunk_size] + data_chunk = self._request.body[data_chunk_start_id:data_chunk_start_id + chunk_size] self._conn.send_data(self.stream_id, data_chunk, end_stream=False) @@ -183,7 +207,7 @@ class Stream: self, self.content_length - self.remaining_content_length, self.content_length, data_frames_sent, - self._metadata['ip_address']) + self._conn_metadata['ip_address']) ) # End the stream if no more data needs to be send @@ -197,24 +221,40 @@ class Stream: # Q. What about the rest of the data? # Ans: Remaining Data frames will be sent when we get a WindowUpdate frame - def receive_window_update(self, delta): + def receive_window_update(self): """Flow control window size was changed. Send data that earlier could not be sent as we were blocked behind the flow control. - - Arguments: - delta -- Window change delta """ - if self.remaining_content_length > 0 and not self.stream_closed_server: + if self.remaining_content_length and not self.stream_closed_server and self.request_sent: self.send_data() def receive_data(self, data: bytes, flow_controlled_length: int): - self._response['body'] += data - self._response_flow_controlled_size += flow_controlled_length + self._response['body'].write(data) + self._response['flow_controlled_size'] += flow_controlled_length + + if self._download_maxsize and self._response['flow_controlled_size'] > self._download_maxsize: + # Clear buffer earlier to avoid keeping data in memory for a long time + self._response['body'].truncate(0) + self.reset_stream(StreamCloseReason.MAXSIZE_EXCEEDED) + return + + if self._download_warnsize \ + and self._response['flow_controlled_size'] > self._download_warnsize \ + and not self._reached_warnsize: + self._reached_warnsize = True + warning_msg = ('Received more ({bytes}) bytes than download ', + 'warn size ({warnsize}) in request {request}') + warning_args = { + 'bytes': self._response['flow_controlled_size'], + 'warnsize': self._download_warnsize, + 'request': self._request + } + LOGGER.warning(warning_msg, warning_args) # Acknowledge the data received self._conn.acknowledge_received_data( - self._response_flow_controlled_size, + self._response['flow_controlled_size'], self.stream_id ) @@ -222,35 +262,83 @@ class Stream: for name, value in headers: self._response['headers'][name] = value - def close(self, event=None): - """Based on the event sent we will handle each case. + # Check if we exceed the allowed max data size which can be received + expected_size = int(self._response['headers'].get(b'Content-Length', -1)) + if self._download_maxsize and expected_size > self._download_maxsize: + self.reset_stream(StreamCloseReason.MAXSIZE_EXCEEDED) + return - event: StreamEnded - Stream is ended by the server hence no further - data or headers should be expected on this stream. - We will call the response deferred callback passing - the response object + if self._download_warnsize and expected_size > self._download_warnsize: + warning_msg = ("Expected response size ({size}) larger than ", + "download warn size ({warnsize}) in request {request}.") + warning_args = { + 'size': expected_size, 'warnsize': self._download_warnsize, + 'request': self._request + } + LOGGER.warning(warning_msg, warning_args) - event: StreamReset - Stream reset via RST_FRAME by the upstream hence forcefully close - this stream and send TODO: ? + def reset_stream(self, reason=StreamCloseReason.RESET): + """Close this stream by sending a RST_FRAME to the remote peer""" + # TODO: Q. REFUSED_STREAM or CANCEL ? + if self.stream_closed_local: + raise StreamClosedError(self.stream_id) - event: None - No event is launched -- Hence we will simply close this stream + self.stream_closed_local = True + self._conn.reset_stream(self.stream_id, ErrorCodes.REFUSED_STREAM) + self._write_to_transport() + self.close(reason) + + def _is_data_lost(self) -> bool: + assert self.stream_closed_server + + expected_size = self._response['flow_controlled_size'] + received_body_size = int(self._response['headers'][b'Content-Length']) + + return expected_size != received_body_size + + def close(self, reason: StreamCloseReason): + """Based on the reason sent we will handle each case. """ # TODO: In case of abruptly stream close # Q1. Do we need to send the request again? # Q2. What response should we send now? - assert self.stream_closed_server is False + if self.stream_closed_server: + raise StreamClosedError(self.stream_id) + self.stream_closed_server = True - if not isinstance(event, StreamEnded): - # TODO - # Stream was abruptly ended here - # Partial - Content-Length header not provided - pass + flags = None + if b'Content-Length' not in self._response['headers']: + # Missing Content-Length - PotentialDataLoss + flags = ['partial'] + elif self._is_data_lost(): + if self._fail_on_dataloss: + self._deferred_response.errback(ResponseFailed([Failure()])) + self._cb_close(self.stream_id) + return + else: + flags = ['dataloss'] + + if reason is StreamCloseReason.ENDED: + self._fire_response_deferred(flags) + + elif reason in (StreamCloseReason.RESET | StreamCloseReason.CONNECTION_LOST): + # Stream was abruptly ended here + self._deferred_response.errback(ResponseFailed([Failure()])) + + elif reason is StreamCloseReason.MAXSIZE_EXCEEDED: + expected_size = int(self._response['headers'].get(b'Content-Length', -1)) + error_msg = ("Cancelling download of {url}: expected response " + "size ({size}) larger than download max size ({maxsize}).") + error_args = { + 'url': self._request.url, + 'size': expected_size, + 'maxsize': self._download_maxsize + } + + LOGGER.error(error_msg, error_args) + self._deferred_response.errback(CancelledError(error_msg.format(**error_args))) - self._fire_response_deferred() self._cb_close(self.stream_id) def _fire_response_deferred(self, flags=None): @@ -258,7 +346,6 @@ class Stream: and fires the response deferred callback with the generated response instance""" # TODO: - # 1. Set flags, certificate, ip_address in response # 2. Should we fire this in case of # 2.1 StreamReset in between when data is received partially # 2.2 Forcefully closed the stream @@ -270,7 +357,7 @@ class Stream: body=self._response['body'] ) - # If there is :status in headers then + # If there is no :status in headers then # HTTP Status Code: 499 - Client Closed Request status = self._response['headers'].get(':status', '499') @@ -278,11 +365,11 @@ class Stream: url=self._request.url, status=status, headers=self._response['headers'], - body=self._response['body'], + body=self._response['body'].getvalue(), request=self._request, flags=flags, - certificate=self._metadata['certificate'], - ip_address=self._metadata['ip_address'] + certificate=self._conn_metadata['certificate'], + ip_address=self._conn_metadata['ip_address'] ) self._deferred_response.callback(response)