@@ -970,45 +970,24 @@ ssize_t Http2Session::OnMaxFrameSizePadding(size_t frameLen,
970970// quite expensive. This is a potential performance optimization target later.
971971void Http2Session::ConsumeHTTP2Data () {
972972 CHECK_NOT_NULL (stream_buf_.base );
973- CHECK_LE (stream_buf_offset_, stream_buf_.len );
974- size_t read_len = stream_buf_.len - stream_buf_offset_;
975973
976974 // multiple side effects.
977- Debug (this , " receiving %d bytes [wants data? %d]" ,
978- read_len,
975+ Debug (this ,
976+ " receiving %d bytes [wants data? %d]" ,
977+ stream_buf_.len ,
979978 nghttp2_session_want_read (session_.get ()));
980- set_receive_paused (false );
981979 custom_recv_error_code_ = nullptr ;
982980 set_receiving ();
983981 ssize_t ret =
984- nghttp2_session_mem_recv (session_.get (),
985- reinterpret_cast <uint8_t *>(stream_buf_.base ) +
986- stream_buf_offset_,
987- read_len);
982+ nghttp2_session_mem_recv (session_.get (),
983+ reinterpret_cast <uint8_t *>(stream_buf_.base ),
984+ stream_buf_.len );
988985 set_receiving (false );
989986 CHECK_NE (ret, NGHTTP2_ERR_NOMEM );
990987 CHECK_IMPLIES (custom_recv_error_code_ != nullptr , ret < 0 );
991988
992- if (is_receive_paused ()) {
993- CHECK (is_reading_stopped ());
994-
995- CHECK_GT (ret, 0 );
996- CHECK_LE (static_cast <size_t >(ret), read_len);
997-
998- // Mark the remainder of the data as available for later consumption.
999- // Even if all bytes were received, a paused stream may delay the
1000- // nghttp2_on_frame_recv_callback which may have an END_STREAM flag.
1001- stream_buf_offset_ += ret;
1002- // Still complete a Close() deferred during mem_recv; do not fall through
1003- // to SendPendingData() here (paused receives historically skip that flush
1004- // because a write may already be in progress).
1005- MaybeFinishPendingClose ();
1006- goto done;
1007- }
1008-
1009989 // We are done processing the current input chunk.
1010990 DecrementCurrentSessionMemory (stream_buf_.len );
1011- stream_buf_offset_ = 0 ;
1012991 stream_buf_ab_.Reset ();
1013992 stream_buf_allocation_.reset ();
1014993 stream_buf_ = uv_buf_init (nullptr , 0 );
@@ -1017,14 +996,6 @@ void Http2Session::ConsumeHTTP2Data() {
1017996 // not written after pending RST_STREAM frames.
1018997 MaybeFinishPendingClose ();
1019998
1020- done:
1021- // Finish a Close() deferred above before flushing, so GOAWAY is not written
1022- // after pending RST_STREAM frames.
1023- if (is_close_pending () && !is_destroyed ()) {
1024- set_close_pending (false );
1025- FinishClose (pending_close_code_, pending_close_socket_closed_);
1026- }
1027-
1028999 // Send any data that was queued up while processing the received data.
10291000 if (ret >= 0 && !is_destroyed ()) {
10301001 SendPendingData ();
@@ -1478,15 +1449,6 @@ int Http2Session::OnDataChunkReceived(nghttp2_session* handle,
14781449 }
14791450 } while (len != 0 );
14801451
1481- // If we are currently waiting for a write operation to finish, we should
1482- // tell nghttp2 that we want to wait before we process more input data.
1483- if (session->is_write_in_progress ()) {
1484- CHECK (session->is_reading_stopped ());
1485- session->set_receive_paused ();
1486- Debug (session, " receive paused" );
1487- return NGHTTP2_ERR_PAUSE ;
1488- }
1489-
14901452 return 0 ;
14911453}
14921454
@@ -1569,7 +1531,6 @@ void Http2StreamListener::OnStreamRead(ssize_t nread, const uv_buf_t& buf) {
15691531 size_t offset = buf.base - session->stream_buf_ .base ;
15701532
15711533 // Verify that the data offset is inside the current read buffer.
1572- CHECK_GE (offset, session->stream_buf_offset_ );
15731534 CHECK_LE (offset, session->stream_buf_ .len );
15741535 CHECK_LE (offset + buf.len , session->stream_buf_ .len );
15751536
@@ -1883,11 +1844,6 @@ void Http2Session::OnStreamAfterWrite(WriteWrap* w, int status) {
18831844 return ;
18841845 }
18851846
1886- // If there is more incoming data queued up, consume it.
1887- if (stream_buf_offset_ > 0 ) {
1888- ConsumeHTTP2Data ();
1889- }
1890-
18911847 if (!is_write_scheduled () && !is_destroyed ()) {
18921848 // Schedule a new write if nghttp2 wants to send data.
18931849 MaybeScheduleWrite ();
@@ -1934,7 +1890,7 @@ void Http2Session::MaybeStopReading() {
19341890 if (is_reading_stopped () || is_closing ()) return ;
19351891 int want_read = nghttp2_session_want_read (session_.get ());
19361892 Debug (this , " wants read? %d" , want_read);
1937- if (want_read == 0 || is_write_in_progress () ) {
1893+ if (want_read == 0 ) {
19381894 set_reading_stopped ();
19391895 stream_->ReadStop ();
19401896 }
@@ -2199,7 +2155,7 @@ void Http2Session::OnStreamRead(ssize_t nread, const uv_buf_t& buf_) {
21992155 Context::Scope context_scope (env ()->context ());
22002156 Http2Scope h2scope (this );
22012157 CHECK_NOT_NULL (stream_);
2202- Debug (this , " receiving %d bytes, offset %d " , nread, stream_buf_offset_ );
2158+ Debug (this , " receiving %d bytes" , nread);
22032159 std::unique_ptr<BackingStore> bs = env ()->release_managed_buffer (buf_);
22042160
22052161 // Only pass data on if nread > 0
@@ -2214,40 +2170,18 @@ void Http2Session::OnStreamRead(ssize_t nread, const uv_buf_t& buf_) {
22142170
22152171 statistics_.data_received += nread;
22162172
2217- if (stream_buf_offset_ == 0 && static_cast <size_t >(nread) != bs->ByteLength ())
2218- [[likely]] {
2173+ // ConsumeHTTP2Data() always consumes the whole chunk, so there is never a
2174+ // partially processed buffer left over from a previous read.
2175+ DCHECK_NULL (stream_buf_.base );
2176+
2177+ if (static_cast <size_t >(nread) != bs->ByteLength ()) [[likely]] {
22192178 // Shrink to the actual amount of used data.
22202179 std::unique_ptr<BackingStore> old_bs = std::move (bs);
22212180 bs = ArrayBuffer::NewBackingStore (
22222181 env ()->isolate (),
22232182 nread,
22242183 BackingStoreInitializationMode::kUninitialized );
22252184 memcpy (bs->Data (), old_bs->Data (), nread);
2226- } else {
2227- // This is a very unlikely case, and should only happen if the ReadStart()
2228- // call in OnStreamAfterWrite() immediately provides data. If that does
2229- // happen, we concatenate the data we received with the already-stored
2230- // pending input data, slicing off the already processed part.
2231- size_t pending_len = stream_buf_.len - stream_buf_offset_;
2232- std::unique_ptr<BackingStore> new_bs = ArrayBuffer::NewBackingStore (
2233- env ()->isolate (),
2234- pending_len + nread,
2235- BackingStoreInitializationMode::kUninitialized );
2236- memcpy (static_cast <char *>(new_bs->Data ()),
2237- stream_buf_.base + stream_buf_offset_,
2238- pending_len);
2239- memcpy (static_cast <char *>(new_bs->Data ()) + pending_len,
2240- bs->Data (),
2241- nread);
2242-
2243- bs = std::move (new_bs);
2244- nread = bs->ByteLength ();
2245- stream_buf_offset_ = 0 ;
2246- stream_buf_ab_.Reset ();
2247-
2248- // We have now fully processed the stream_buf_ input chunk (by moving the
2249- // remaining part into buf, which will be accounted for below).
2250- DecrementCurrentSessionMemory (stream_buf_.len );
22512185 }
22522186
22532187 IncrementCurrentSessionMemory (nread);
0 commit comments