From 089dbc75e78a2da9c455f21bb3c7ebaaeb2e3582 Mon Sep 17 00:00:00 2001 From: Aditya Date: Wed, 17 Jun 2020 20:57:03 +0530 Subject: [PATCH] chore: use deque for pending request pool - Use itertools.count to generate next stream_id BREAKING CHANGES When sending data/body more than the local flow control window -- no window update occurs to send the remaining data frames. Hence, the complete body is not send resulting in no response received. --- scrapy/core/http2/protocol.py | 16 +++++++++------- scrapy/core/http2/stream.py | 4 ++-- 2 files changed, 11 insertions(+), 9 deletions(-) diff --git a/scrapy/core/http2/protocol.py b/scrapy/core/http2/protocol.py index 915284c33..dbc048ffa 100644 --- a/scrapy/core/http2/protocol.py +++ b/scrapy/core/http2/protocol.py @@ -1,4 +1,6 @@ +import itertools import logging +from collections import deque from h2.config import H2Configuration from h2.connection import H2Connection @@ -35,7 +37,7 @@ class H2ClientProtocol(Protocol): # ID of the next request stream # Following the convention made by hyper-h2 each client ID # will be odd. - self.next_stream_id = 1 + self.stream_id_count = itertools.count(start=1, step=2) # Streams are stored in a dictionary keyed off their stream IDs self.streams = {} @@ -45,7 +47,7 @@ class H2ClientProtocol(Protocol): # we keep all requests in a pool and send them as the connection # is made self.is_connection_made = False - self._pending_request_stream_pool = [] + self._pending_request_stream_pool = deque() def _stream_close_cb(self, stream_id: int): """Called when stream is closed completely @@ -55,8 +57,7 @@ class H2ClientProtocol(Protocol): def _new_stream(self, request: Request): """Instantiates a new Stream object """ - stream_id = self.next_stream_id - self.next_stream_id += 2 + stream_id = next(self.stream_id_count) stream = Stream( stream_id=stream_id, @@ -72,11 +73,10 @@ class H2ClientProtocol(Protocol): def _send_pending_requests(self): # TODO: handle MAX_CONCURRENT_STREAMS # Initiate all pending requests - for stream in self._pending_request_stream_pool: + while len(self._pending_request_stream_pool): + stream = self._pending_request_stream_pool.popleft() stream.initiate_request() - self._pending_request_stream_pool.clear() - def _write_to_transport(self): """ Write data to the underlying transport connection from the HTTP2 connection instance if any @@ -108,6 +108,8 @@ class H2ClientProtocol(Protocol): self._write_to_transport() self.is_connection_made = True + # Send off all the pending requests + # as now we have established a proper HTTP/2 connection self._send_pending_requests() def dataReceived(self, data): diff --git a/scrapy/core/http2/stream.py b/scrapy/core/http2/stream.py index a0b75850d..c2e1adce5 100644 --- a/scrapy/core/http2/stream.py +++ b/scrapy/core/http2/stream.py @@ -132,7 +132,7 @@ class Stream: bytes_to_send_size = min(window_size, self.remaining_content_length) # We now need to send a number of data frames. - while bytes_to_send_size > 0: + while bytes_to_send_size: chunk_size = min(bytes_to_send_size, max_frame_size) data_chunk_start_id = self.content_length - self.remaining_content_length @@ -163,7 +163,7 @@ class Stream: Arguments: delta -- Window change delta """ - if self.remaining_content_length > 0 and not self.stream_closed_local: + if self.stream_closed_local is False: self.send_data() def receive_data(self, data: bytes, flow_controlled_length: int):