This is an automated email from the ASF dual-hosted git repository.
pjfanning pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/pekko-http.git
The following commit(s) were added to refs/heads/main by this push:
new 2dc91fef9 fix: release connection-level flow control accounting for
discarded buffered data (#1259)
2dc91fef9 is described below
commit 2dc91fef995e89564186fcdf927dccde2193f462
Author: PJ Fanning <[email protected]>
AuthorDate: Sun Sep 6 12:05:28 2026 +0100
fix: release connection-level flow control accounting for discarded
buffered data (#1259)
Motivation:
`totalBufferedData` feeds the connection-level flow controller, which
emits a WINDOW_UPDATE only while `outstanding + buffered` stays below
half of `incoming-connection-level-buffer-size`. It is incremented for
every DATA frame received but was not decremented on the three paths
that discard a non-empty buffer:
- `IncomingStreamBuffer.onRstStreamFrame` cleared the buffer after the
peer reset the stream,
- `IncomingStreamBuffer.onDownstreamFinish` cleared it when the
application cancelled the entity stream,
- `CollectingIncomingData.onRstStreamFrame` was a no-op even though the
bytes collected so far had been counted.
A peer that sends data the handler does not read and then resets the
stream therefore ratchets `totalBufferedData` up permanently. Once the
leaked total reaches half the configured buffer size the server stops
replenishing the connection window, it drains to zero and every stream
on that connection stalls.
Modification:
Add `IncomingStreamBuffer.discardBuffer()`, which subtracts what is
still buffered from `totalBufferedData` before clearing it, and use it
on both discard paths. Subtract the collected bytes in
`CollectingIncomingData.onRstStreamFrame` for the same reason.
Result:
Resetting or cancelling a stream releases the connection-level window
its buffered data reserved, so the connection keeps being replenished
and no longer stalls.
Tests:
- sbt "http2-tests/testOnly
org.apache.pekko.http.impl.engine.http2.Http2ServerSpec
org.apache.pekko.http.impl.engine.http2.Http2ClientSpec
org.apache.pekko.http.impl.engine.http2.Http2ClientServerSpec" - pass (174
tests); a new test buffers request data the handler never reads and resets the
stream six times over, with the connection-level buffer size lowered to the
initial window so the effect is reached quickly, then asserts the connection
still accepts a request. Verified it fails [...]
- sbt http-core/mimaReportBinaryIssues - pass (internal impl.engine.http2
change, no public API).
References:
None - releases buffered-data accounting when a stream is reset or cancelled
---
.../impl/engine/http2/Http2StreamHandling.scala | 18 +++++++++++++++---
.../http/impl/engine/http2/Http2ServerSpec.scala | 22 ++++++++++++++++++++++
2 files changed, 37 insertions(+), 3 deletions(-)
diff --git
a/http-core/src/main/scala/org/apache/pekko/http/impl/engine/http2/Http2StreamHandling.scala
b/http-core/src/main/scala/org/apache/pekko/http/impl/engine/http2/Http2StreamHandling.scala
index 5de3aa173..364e523fc 100644
---
a/http-core/src/main/scala/org/apache/pekko/http/impl/engine/http2/Http2StreamHandling.scala
+++
b/http-core/src/main/scala/org/apache/pekko/http/impl/engine/http2/Http2StreamHandling.scala
@@ -391,7 +391,9 @@ private[http2] trait Http2StreamHandling extends
GraphStageLogic with LogHelper
override protected def onTrailer(parsedHeadersFrame: ParsedHeadersFrame):
StreamState = this // trailing headers not supported for requests right now
override protected def incrementWindow(delta: Int): StreamState =
copy(extraInitialWindow = extraInitialWindow + delta)
- override protected def onRstStreamFrame(rstStreamFrame: RstStreamFrame):
Unit = {} // nothing to do here
+ // the data collected so far is dropped together with this state, so stop
reserving connection-level window for it
+ override protected def onRstStreamFrame(rstStreamFrame: RstStreamFrame):
Unit =
+ totalBufferedData -= collectedData.length
}
case class OpenReceivingDataFirst(buffer: IncomingStreamBuffer,
extraInitialWindow: Int = 0)
extends
ReceivingDataWithBuffer(HalfClosedRemoteWaitingForOutgoingStream(extraInitialWindow))
{
@@ -640,10 +642,20 @@ private[http2] trait Http2StreamHandling extends
GraphStageLogic with LogHelper
streamStates.remove(streamId)
headRequestStreamIds -= streamId
wasClosed = true
- buffer = ByteString.empty
+ discardBuffer()
trailingHeaders = None
}
+ /**
+ * Drops what is still buffered and releases the connection-level window
it reserved. Without the release the
+ * connection-level flow controller keeps counting these bytes as buffered
forever, so it stops replenishing the
+ * connection window and every stream on the connection eventually stalls.
+ */
+ private def discardBuffer(): Unit = {
+ totalBufferedData -= buffer.length
+ buffer = ByteString.empty
+ }
+
def isDone: Boolean = outlet.isClosed
def onDataFrame(data: DataFrame): Unit =
@@ -676,7 +688,7 @@ private[http2] trait Http2StreamHandling extends
GraphStageLogic with LogHelper
}
def onRstStreamFrame(rst: RstStreamFrame): Unit = {
outlet.fail(new PeerClosedStreamException(rst.streamId, rst.errorCode))
- buffer = ByteString.empty
+ discardBuffer()
trailingHeaders = None
wasClosed = true
}
diff --git
a/http2-tests/src/test/scala/org/apache/pekko/http/impl/engine/http2/Http2ServerSpec.scala
b/http2-tests/src/test/scala/org/apache/pekko/http/impl/engine/http2/Http2ServerSpec.scala
index 78bd9d78b..c90661b85 100644
---
a/http2-tests/src/test/scala/org/apache/pekko/http/impl/engine/http2/Http2ServerSpec.scala
+++
b/http2-tests/src/test/scala/org/apache/pekko/http/impl/engine/http2/Http2ServerSpec.scala
@@ -346,6 +346,28 @@ class Http2ServerSpec extends Http2SpecWithMaterializer("""
trailingResponseHeaders should contain(("x-good", "fine"))
trailingResponseHeaders.map(_._1) should not contain "x-evil"
})
+ "release connection-level flow control accounting when a stream with
buffered data is reset"
+ .inAssertAllStagesStopped(new TestSetup with RequestResponseProbes {
+ override def settings: ServerSettings =
+
super.settings.mapHttp2Settings(_.withIncomingConnectionLevelBufferSize(Http2Protocol.InitialWindowSize))
+
+ // Each round buffers request data that the handler never reads and
then resets the stream. Those bytes must
+ // be released from the connection-level accounting: otherwise the
flow controller keeps counting them as
+ // buffered forever, stops replenishing the connection window, and
the peer runs out of window entirely.
+ (0 until 6).foreach { i =>
+ val streamId = 1 + i * 2
+ network.sendHEADERS(streamId, endStream = false,
+ network.headersForRequest(HttpRequest(HttpMethods.POST, "/")))
+ user.expectRequest() // the entity is deliberately never read, so
the data stays buffered
+ network.sendDATA(streamId, endStream = false, ByteString(new
Array[Byte](20000)))
+ network.sendRST_STREAM(streamId, ErrorCode.CANCEL)
+ network.pollForWindowUpdates(100.millis)
+ }
+
+ // the connection is still usable because the window was replenished
along the way
+ network.sendHEADERS(13, endStream = true,
network.headersForRequest(Get("/")))
+ user.expectRequest()
+ })
"consider stream as closed after sending out strict response >
WINDOW_SIZE".inAssertAllStagesStopped(
new TestSetup with RequestResponseProbes {
override def settings: ServerSettings =
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]