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 478c58bb7 feat: add max-connection-age setting for HTTP/2 server
connections (#1316)
478c58bb7 is described below
commit 478c58bb733fadbe921660bccc964311ad58e9ef
Author: Anders Kreinøe <[email protected]>
AuthorDate: Fri Oct 2 11:20:19 2026 +0200
feat: add max-connection-age setting for HTTP/2 server connections (#1316)
* feat: add max-connection-age setting for HTTP/2 server connections
Motivation:
Long-lived HTTP/2 connections (as used by gRPC) lead to an uneven load
distribution across server instances: clients stay connected to the
instances they found at connect time, and instances added later (after
a scale-out or a rolling deploy) receive no share of the existing
traffic. The server is the side that can retire a connection
gracefully, via GOAWAY. grpc-java offers this as maxConnectionAge;
pekko-http has no equivalent (akka/akka-grpc#967 is the corresponding
request on the Akka side).
Modification:
Add a `pekko.http.server.http2.max-connection-age` setting, default
`infinite` (disabled). When a server connection reaches the configured
age, the existing graceful termination path is triggered:
GOAWAY(NO_ERROR) is sent, streams that are in flight complete normally,
streams opened after the GOAWAY are refused with
RST_STREAM(REFUSED_STREAM), and the connection is closed once no
streams remain. The age is jittered per connection by a configurable
fraction, `max-connection-age-jitter` (default 0.1 = +/- 10%, the value
grpc-java applies; 0 disables jitter), so that connections that were
opened together are not all closed at the same time.
`triggerTermination` now accepts an infinite deadline, in which case no
forced-close timer is scheduled. Also corrects the termination debug
log, which printed the timer key instead of the deadline.
Result:
Operators can cap the lifetime of server-side HTTP/2 connections to
rebalance long-lived connections across server instances. Behavior is
unchanged by default.
Tests:
- sbt "http2-tests/test": 376 tests pass, including 2 new directional
tests for max-connection-age in Http2ServerSpec
- sbt validatePullRequest: passes
- sbt "+http-core/mimaReportBinaryIssues": clean, with 5 new
ReversedMissingMethodProblem filters for the added methods
- sbt checkCodeStyle: clean; sbt headerCreateAll: no changes
- sbt docs/paradox: builds, remaining warnings pre-existing
- sbt sortImports: environment failure unrelated to this change
(http/scalafixAll fails with NoSuchMethodError in scala.meta on a
clean checkout of main too); the files changed here are sort-clean
- manual end-to-end check against grpc-java 1.75.0: with
max-connection-age = 5s, a client calling every 200 ms for 16 s saw
0 failures across 3 connection retirements, and a unary call that
was in flight when the GOAWAY was sent completed normally
References:
None - no pekko-http issue tracks this; akka/akka-grpc#967 is the
equivalent request against akka-http/akka-grpc.
* Set @since of the new settings to 1.4.2 as requested in review
* Revert "Set @since of the new settings to 1.4.2 as requested in review"
This reverts commit 9ce450d826db93c02a9089d8a882d449456b25ee.
* Map ChronoUnit.FOREVER back to Duration.Inf in the Java
withMaxConnectionAge
withMaxConnectionAge(getMaxConnectionAge) threw ArithmeticException on the
default settings: toMillis overflows on ChronoUnit.FOREVER.getDuration, the
value getMaxConnectionAge returns for an infinite age. Adds the inverse of
JavaDurationConverter.toJava and a round-trip test.
* Document when JavaDurationConverter.toScala throws
* Make the infinite round-trip test independent of the default
max-connection-age
* Add max-connection-age-grace to bound the drain after max-connection-age
The drain started by the MaxConnectionAge timer had no deadline, so a
stream that never completes kept the connection open indefinitely. The
new setting bounds it, with a default of 30s. The value infinite keeps
the previous behaviour of waiting for all requests in flight to complete.
* Drop the grpc-java aside from the max-connection-age-grace comment
* Let a later terminate shorten a termination that is already in progress
triggerTermination ignored every call after the first, so a server binding
termination could not enforce its deadline on a connection that was already
draining after max-connection-age. Now an earlier deadline reschedules the
forced close, a later one is ignored. The max-connection-age timer is
cancelled once any termination starts, so the age and its grace period never
shorten a termination in progress.
* Guard the max-connection-age timer instead of cancelling it when a
termination starts
Same behaviour, but the rule is visible where it applies: the timer handler
checks whether a termination is already in progress and logs that there is
nothing to do, rather than relying on the cancelled timer's message being
dropped by the stage.
* Name scheduleForcedCloseIfEarlier after what it does instead of
commenting it
* Say what an age expiry does on a connection that is already terminating,
next to the age setting
* State that max-connection-age applies to HTTP/2 only, and reword the
jitter test comment
* Do not point server-side forced-close logging at the client-only
completion-timeout setting
The message is shared by both sides. On the client the deadline is the
completion-timeout setting, on the server it is the deadline passed to
terminate or the max-connection-age-grace setting.
* Do not blame the peer in the server-side forced-close log message
---
docs/src/main/paradox/server-side/http2.md | 38 +++++++
.../http2-max-connection-age.excludes | 25 +++++
http-core/src/main/resources/reference.conf | 31 ++++++
.../pekko/http/impl/engine/http2/Http2Demux.scala | 70 ++++++++++--
.../http/impl/util/JavaDurationConverter.scala | 11 ++
.../javadsl/settings/Http2ServerSettings.scala | 55 ++++++++++
.../scaladsl/settings/Http2ServerSettings.scala | 57 ++++++++++
.../settings/Http2ServerSettingsSpec.scala | 69 ++++++++++++
.../http/impl/engine/http2/Http2ServerSpec.scala | 119 +++++++++++++++++++++
9 files changed, 467 insertions(+), 8 deletions(-)
diff --git a/docs/src/main/paradox/server-side/http2.md
b/docs/src/main/paradox/server-side/http2.md
index f63a89bcc..a34011483 100644
--- a/docs/src/main/paradox/server-side/http2.md
+++ b/docs/src/main/paradox/server-side/http2.md
@@ -74,6 +74,44 @@ server supports HTTP/2.
For this reason the approach is known as HTTP/2 with
[Prior Knowledge](https://www.rfc-editor.org/rfc/rfc9113.html#section-3.3).
+## Limiting the lifetime of connections
+
+Long-lived HTTP/2 connections, for example those used by gRPC, can lead to an
uneven load distribution across
+server instances: clients stay connected to the instances they found at
connect time, and instances added later
+(after a scale-out or a rolling deploy) receive no share of the existing
traffic.
+
+To rebalance connections regularly, set a maximum connection age:
+
+```
+pekko.http.server.http2.max-connection-age = 120s
+```
+
+When the age of a connection exceeds the configured value, the server sends a
GOAWAY frame, lets requests that
+are already in flight complete, and then closes the connection. Streams that
the peer opens after the GOAWAY
+frame was sent are refused, upon which well-behaved clients (for example,
grpc-java) transparently retry them on
+a new connection.
+
+The setting applies to HTTP/2 connections only: HTTP/1.1 connections accepted
on the same port are not
+age-limited.
+
+If the connection is already being terminated when its age expires, for
example because the server binding is
+being terminated, the expiry has no effect and the termination in progress
keeps its own deadline. Conversely,
+terminating the server binding with an earlier deadline shortens a drain that
the age started.
+
+Requests that are still in flight when the age expires are given a grace
period to complete, after which the
+connection is closed even if they have not completed:
+
+```
+pekko.http.server.http2.max-connection-age-grace = 30s
+```
+
+The default grace period is 30 seconds. Set it to `infinite` to wait for all
requests in flight to complete,
+however long they take.
+
+A jitter is applied to the configured value for each connection (by default
+/- 10%, configurable via
+`pekko.http.server.http2.max-connection-age-jitter`), so that connections that
were opened together are not
+all closed at the same time.
+
## Trailing headers
Like in the [HTTP/1.1 'Chunked' transfer
encoding](https://datatracker.ietf.org/doc/html/rfc7230#section-4.1.2),
diff --git
a/http-core/src/main/mima-filters/2.0.x.backwards.excludes/http2-max-connection-age.excludes
b/http-core/src/main/mima-filters/2.0.x.backwards.excludes/http2-max-connection-age.excludes
new file mode 100644
index 000000000..a03a9eac8
--- /dev/null
+++
b/http-core/src/main/mima-filters/2.0.x.backwards.excludes/http2-max-connection-age.excludes
@@ -0,0 +1,25 @@
+# Licensed to the Apache Software Foundation (ASF) under one
+# or more contributor license agreements. See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership. The ASF licenses this file
+# to you under the Apache License, Version 2.0 (the
+# "License"); you may not use this file except in compliance
+# with the License. You may obtain a copy of the License at
+#
+# http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing,
+# software distributed under the License is distributed on an
+# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+# KIND, either express or implied. See the License for the
+# specific language governing permissions and limitations
+# under the License.
+
+# new max-connection-age settings for HTTP/2 servers
+ProblemFilters.exclude[ReversedMissingMethodProblem]("org.apache.pekko.http.scaladsl.settings.Http2ServerSettings.maxConnectionAge")
+ProblemFilters.exclude[ReversedMissingMethodProblem]("org.apache.pekko.http.scaladsl.settings.Http2ServerSettings.maxConnectionAgeGrace")
+ProblemFilters.exclude[ReversedMissingMethodProblem]("org.apache.pekko.http.scaladsl.settings.Http2ServerSettings.maxConnectionAgeJitter")
+ProblemFilters.exclude[ReversedMissingMethodProblem]("org.apache.pekko.http.javadsl.settings.Http2ServerSettings.withMaxConnectionAgeJitter")
+ProblemFilters.exclude[ReversedMissingMethodProblem]("org.apache.pekko.http.impl.engine.http2.Http2Demux.maxConnectionAge")
+ProblemFilters.exclude[ReversedMissingMethodProblem]("org.apache.pekko.http.impl.engine.http2.Http2Demux.maxConnectionAgeGrace")
+ProblemFilters.exclude[ReversedMissingMethodProblem]("org.apache.pekko.http.impl.engine.http2.Http2Demux.maxConnectionAgeJitter")
diff --git a/http-core/src/main/resources/reference.conf
b/http-core/src/main/resources/reference.conf
index 8ca31252f..bea6d5582 100644
--- a/http-core/src/main/resources/reference.conf
+++ b/http-core/src/main/resources/reference.conf
@@ -352,6 +352,37 @@ pekko.http {
# When zero the ping-interval is used, if set the value must be evenly
divisible by less than or equal to the ping-interval.
ping-timeout = 0s
+ # The maximum time a connection is kept open before the server closes it
gracefully. When the age of a
+ # connection exceeds this value, the server sends a GOAWAY frame, lets
requests that are already in flight
+ # complete (see `max-connection-age-grace`), and then closes the
connection. Streams that the peer opens
+ # after the GOAWAY frame was sent are refused with a
RST_STREAM(REFUSED_STREAM) frame, upon which
+ # well-behaved clients retry them on a new connection.
+ #
+ # Closing connections regularly helps to rebalance long-lived HTTP/2
connections (as used by gRPC) across
+ # server instances, for example after a scale-out or a rolling deploy.
+ #
+ # The setting applies to HTTP/2 connections only: HTTP/1.1 connections
accepted on the same port are not
+ # age-limited.
+ #
+ # If the connection is already being terminated when its age expires,
for example because the server
+ # binding is being terminated, the expiry has no effect and the
termination in progress keeps its own
+ # deadline.
+ #
+ # The value `infinite` disables this mechanism and is the default.
+ max-connection-age = infinite
+
+ # The time that requests in flight are given to complete after a
connection reached `max-connection-age`
+ # and the GOAWAY frame was sent. When the grace period expires, the
connection is closed even if requests
+ # are still in flight. The value `infinite` disables the limit: the
connection is closed only once all
+ # requests in flight have completed.
+ max-connection-age-grace = 30s
+
+ # The jitter applied to `max-connection-age`, as a fraction of the
configured value: with the default of
+ # 0.1 each connection is closed after between 90% and 110% of the
configured age, so that connections that
+ # were opened together are not all closed at the same time (grpc-java
applies the same 10% jitter).
+ # Set to 0 to disable jitter. Must be >= 0 and < 1.
+ max-connection-age-jitter = 0.1
+
frame-type-throttle {
# Configure the throttle for non-data frame types
(https://github.com/apache/pekko-http/issues/332).
# The supported frame-types for throttling are:
diff --git
a/http-core/src/main/scala/org/apache/pekko/http/impl/engine/http2/Http2Demux.scala
b/http-core/src/main/scala/org/apache/pekko/http/impl/engine/http2/Http2Demux.scala
index 35ff328cc..4e3cd4854 100644
---
a/http-core/src/main/scala/org/apache/pekko/http/impl/engine/http2/Http2Demux.scala
+++
b/http-core/src/main/scala/org/apache/pekko/http/impl/engine/http2/Http2Demux.scala
@@ -44,9 +44,13 @@ import pekko.stream.stage.{
import pekko.util.ByteString
import pekko.util.OptionVal
+import java.util.concurrent.ThreadLocalRandom
+
import scala.concurrent.{ ExecutionContext, Future, Promise }
+import scala.concurrent.duration.Deadline
import scala.concurrent.duration.Duration
import scala.concurrent.duration.DurationInt
+import scala.concurrent.duration.DurationLong
import scala.concurrent.duration.FiniteDuration
import scala.util.control.NonFatal
@@ -67,6 +71,11 @@ private[http2] class Http2ClientDemux(http2Settings:
Http2ClientSettings, master
}
override def completionTimeout: FiniteDuration =
http2Settings.completionTimeout
+
+ // a maximum connection age is not supported on the client side
+ def maxConnectionAge: Duration = Duration.Inf
+ def maxConnectionAgeGrace: Duration = Duration.Inf
+ def maxConnectionAgeJitter: Double = 0.0
}
/**
@@ -81,6 +90,10 @@ private[http2] class Http2ServerDemux(http2Settings:
Http2ServerSettings, initia
def completionTimeout: FiniteDuration =
throw new IllegalArgumentException("Completion timeout not supported for
servers")
+
+ def maxConnectionAge: Duration = http2Settings.maxConnectionAge
+ def maxConnectionAgeGrace: Duration = http2Settings.maxConnectionAgeGrace
+ def maxConnectionAgeJitter: Double = http2Settings.maxConnectionAgeJitter
}
/**
@@ -230,13 +243,16 @@ private[http2] abstract class Http2Demux(http2Settings:
Http2CommonSettings,
def wrapTrailingHeaders(headers: ParsedHeadersFrame):
Option[HttpEntity.ChunkStreamPart]
def completionTimeout: FiniteDuration
+ def maxConnectionAge: Duration
+ def maxConnectionAgeGrace: Duration
+ def maxConnectionAgeJitter: Double
override def createLogicAndMaterializedValue(inheritedAttributes:
Attributes): (GraphStageLogic, ServerTerminator) = {
object Logic extends TimerGraphStageLogic(shape) with
Http2MultiplexerSupport with Http2StreamHandling
with GenericOutletSupport with StageLogging with LogHelper with
ServerTerminator {
logic =>
- import Http2Demux.CompletionTimeout
+ import Http2Demux.{ CompletionTimeout, MaxConnectionAge }
def wrapTrailingHeaders(headers: ParsedHeadersFrame):
Option[HttpEntity.ChunkStreamPart] =
stage.wrapTrailingHeaders(headers)
@@ -256,24 +272,38 @@ private[http2] abstract class Http2Demux(http2Settings:
Http2CommonSettings,
private val terminationPromise = Promise[Http.HttpTerminated]()
private var terminating: Boolean = false
private var lastIdBeforeTermination: Int = 0
+ // when the forced close of a terminating connection is due, unset while
no forced close is scheduled
+ private var forcedCloseDeadline: OptionVal[Deadline] = OptionVal.None
private val terminateCallback =
getAsyncCallback[FiniteDuration](triggerTermination)
override def terminate(deadline: FiniteDuration)(implicit ex:
ExecutionContext): Future[Http.HttpTerminated] = {
terminateCallback.invoke(deadline)
terminationPromise.future
}
- private def triggerTermination(deadline: FiniteDuration): Unit =
- // check if we are already terminating, otherwise start termination
+ private def triggerTermination(deadline: Duration): Unit = {
if (!terminating) {
log.debug(
"Termination of this connection was triggered. Sending GOAWAY and
waiting for open requests to complete for {}.",
- CompletionTimeout)
+ deadline)
terminating = true
pushGOAWAY(ErrorCode.NO_ERROR, "Voluntary connection close.")
lastIdBeforeTermination = lastStreamId()
completeIfDone()
- if (!isClosed(frameOut))
- scheduleOnce(CompletionTimeout, deadline)
}
+ scheduleForcedCloseIfEarlier(deadline)
+ }
+ private def scheduleForcedCloseIfEarlier(deadline: Duration): Unit =
deadline match {
+ case deadline: FiniteDuration if !isClosed(frameOut) =>
+ val due = Deadline.now + deadline
+ val earlier = forcedCloseDeadline match {
+ case OptionVal.Some(scheduled) => due < scheduled
+ case _ => true
+ }
+ if (earlier) {
+ forcedCloseDeadline = OptionVal.Some(due)
+ scheduleOnce(CompletionTimeout, deadline)
+ }
+ case _ => // no deadline, wait for open requests to complete
+ }
def frameOutFinished(): Unit = {
// make sure we clean up/fail substreams with a custom failure before
stage is canceled
@@ -315,6 +345,15 @@ private[http2] abstract class Http2Demux(http2Settings:
Http2CommonSettings,
pingState.tickInterval().foreach(interval =>
// to limit overhead rather than constantly rescheduling a timer and
looking at system time we use a constant timer
scheduleAtFixedRate(ConfigurablePing.Tick, interval, interval))
+
+ maxConnectionAge match {
+ case age: FiniteDuration =>
+ // The age of each connection is jittered so that connections that
were opened together are not
+ // all closed at the same time, see `max-connection-age-jitter` in
the configuration.
+ val jitterFactor = 1.0 + maxConnectionAgeJitter * (2 *
ThreadLocalRandom.current().nextDouble() - 1)
+ scheduleOnce(MaxConnectionAge, (age.toMillis *
jitterFactor).toLong.max(1L).millis)
+ case _ => // no maximum connection age configured
+ }
}
override def pushGOAWAY(errorCode: ErrorCode, debug: String): Unit = {
@@ -522,9 +561,23 @@ private[http2] abstract class Http2Demux(http2Settings:
Http2CommonSettings,
} else {
pingState.clear()
}
+ case MaxConnectionAge =>
+ // the max-connection-age and its grace period must not shorten the
deadline of a termination
+ // that is already in progress
+ if (terminating)
+ debug(
+ "Connection reached the configured max-connection-age while a
termination is already in progress, nothing to do")
+ else {
+ debug("Connection reached the configured max-connection-age,
closing it gracefully")
+ triggerTermination(maxConnectionAgeGrace)
+ }
case CompletionTimeout =>
- info(
- "Timeout: Peer didn't finish in-flight requests. Closing pending
HTTP/2 streams. Increase this timeout via the 'completion-timeout' setting.")
+ if (isServer)
+ info(
+ "Timeout: requests in flight did not complete within the
termination deadline (the deadline passed to terminate, or the
'max-connection-age-grace' setting). Closing pending HTTP/2 streams.")
+ else
+ info(
+ "Timeout: Peer didn't finish in-flight requests. Closing pending
HTTP/2 streams. Increase this timeout via the 'completion-timeout' setting.")
shutdownStreamHandling()
completeStage()
@@ -545,4 +598,5 @@ private[http2] abstract class Http2Demux(http2Settings:
Http2CommonSettings,
@InternalApi
private[pekko] object Http2Demux {
case object CompletionTimeout
+ case object MaxConnectionAge
}
diff --git
a/http-core/src/main/scala/org/apache/pekko/http/impl/util/JavaDurationConverter.scala
b/http-core/src/main/scala/org/apache/pekko/http/impl/util/JavaDurationConverter.scala
index ef13bbde3..3d33ca67e 100644
---
a/http-core/src/main/scala/org/apache/pekko/http/impl/util/JavaDurationConverter.scala
+++
b/http-core/src/main/scala/org/apache/pekko/http/impl/util/JavaDurationConverter.scala
@@ -35,4 +35,15 @@ private[http] object JavaDurationConverter {
case scala.concurrent.duration.Duration.MinusInf =>
ChronoUnit.FOREVER.getDuration.negated()
case _ =>
ChronoUnit.FOREVER.getDuration
}
+
+ /**
+ * Inverse of [[toJava]]: `ChronoUnit.FOREVER.getDuration` is mapped back to
`Duration.Inf`,
+ * every other value is converted to a finite duration.
+ *
+ * @throws IllegalArgumentException if the value is not
`ChronoUnit.FOREVER.getDuration` but too large
+ * for a finite Scala duration (about 292
years)
+ */
+ def toScala(d: java.time.Duration): scala.concurrent.duration.Duration =
+ if (d == ChronoUnit.FOREVER.getDuration)
scala.concurrent.duration.Duration.Inf
+ else d.toScala
}
diff --git
a/http-core/src/main/scala/org/apache/pekko/http/javadsl/settings/Http2ServerSettings.scala
b/http-core/src/main/scala/org/apache/pekko/http/javadsl/settings/Http2ServerSettings.scala
index 855533df5..19d5a177d 100644
---
a/http-core/src/main/scala/org/apache/pekko/http/javadsl/settings/Http2ServerSettings.scala
+++
b/http-core/src/main/scala/org/apache/pekko/http/javadsl/settings/Http2ServerSettings.scala
@@ -17,6 +17,7 @@ import java.time.Duration
import org.apache.pekko
import pekko.annotation.DoNotInherit
+import pekko.http.impl.util.JavaDurationConverter
import pekko.http.scaladsl
import com.typesafe.config.Config
@@ -83,6 +84,60 @@ trait Http2ServerSettings {
def getPingTimeout: Duration = Duration.ofMillis(pingTimeout.toMillis)
def withPingTimeout(timeout: Duration): Http2ServerSettings =
withPingTimeout(timeout.toMillis.millis)
+ /**
+ * The maximum time a connection is kept open before the server closes it
gracefully. When the age of a
+ * connection exceeds this value, the server sends a GOAWAY frame, lets
requests that are already in flight
+ * complete within [[getMaxConnectionAgeGrace]], and then closes the
connection. If the connection is
+ * already being terminated when its age expires, for example because the
server binding is being
+ * terminated, the expiry has no effect and the termination in progress
keeps its own deadline. The value
+ * `ChronoUnit.FOREVER.getDuration` represents an infinite age, which
disables this mechanism and is the
+ * default.
+ *
+ * @since 2.0.0
+ */
+ def getMaxConnectionAge: Duration =
JavaDurationConverter.toJava(maxConnectionAge)
+
+ /**
+ * Pass `ChronoUnit.FOREVER.getDuration` to disable the maximum connection
age.
+ *
+ * @since 2.0.0
+ */
+ def withMaxConnectionAge(age: Duration): Http2ServerSettings =
+ withMaxConnectionAge(JavaDurationConverter.toScala(age))
+
+ /**
+ * The time that requests in flight are given to complete after a connection
reached the maximum
+ * connection age and the GOAWAY frame was sent. When the grace period
expires, the connection is closed
+ * even if requests are still in flight. The value
`ChronoUnit.FOREVER.getDuration` represents an infinite
+ * grace period, which disables the limit, so that the connection is closed
only once all requests in
+ * flight have completed.
+ *
+ * @since 2.0.0
+ */
+ def getMaxConnectionAgeGrace: Duration =
JavaDurationConverter.toJava(maxConnectionAgeGrace)
+
+ /**
+ * Pass `ChronoUnit.FOREVER.getDuration` for an infinite grace period.
+ *
+ * @since 2.0.0
+ */
+ def withMaxConnectionAgeGrace(grace: Duration): Http2ServerSettings =
+ withMaxConnectionAgeGrace(JavaDurationConverter.toScala(grace))
+
+ /**
+ * The jitter applied to the maximum connection age per connection, as a
fraction of the configured age:
+ * with the default of 0.1 each connection is closed after between 90% and
110% of the configured age, so
+ * that connections that were opened together are not all closed at the same
time. 0 disables jitter.
+ *
+ * @since 2.0.0
+ */
+ def getMaxConnectionAgeJitter: Double = maxConnectionAgeJitter
+
+ /**
+ * @since 2.0.0
+ */
+ def withMaxConnectionAgeJitter(jitter: Double): Http2ServerSettings
+
def getFrameTypeThrottleFrameTypes(): java.util.Set[String] =
frameTypeThrottleFrameTypes.asJava
def getFrameTypeThrottleCost(): Int = frameTypeThrottleCost
def getFrameTypeThrottleBurst(): Int = frameTypeThrottleBurst
diff --git
a/http-core/src/main/scala/org/apache/pekko/http/scaladsl/settings/Http2ServerSettings.scala
b/http-core/src/main/scala/org/apache/pekko/http/scaladsl/settings/Http2ServerSettings.scala
index a2762156c..7d5b42bd3 100644
---
a/http-core/src/main/scala/org/apache/pekko/http/scaladsl/settings/Http2ServerSettings.scala
+++
b/http-core/src/main/scala/org/apache/pekko/http/scaladsl/settings/Http2ServerSettings.scala
@@ -135,6 +135,53 @@ trait Http2ServerSettings extends
javadsl.settings.Http2ServerSettings with Http
def pingTimeout: FiniteDuration
def withPingTimeout(timeout: FiniteDuration): Http2ServerSettings =
copy(pingTimeout = timeout)
+ /**
+ * The maximum time a connection is kept open before the server closes it
gracefully. When the age of a
+ * connection exceeds this value, the server sends a GOAWAY frame, lets
requests that are already in flight
+ * complete within [[maxConnectionAgeGrace]], and then closes the
connection. If the connection is already
+ * being terminated when its age expires, for example because the server
binding is being terminated, the
+ * expiry has no effect and the termination in progress keeps its own
deadline. The value `Duration.Inf`
+ * disables this mechanism and is the default.
+ *
+ * @since 2.0.0
+ */
+ def maxConnectionAge: Duration
+
+ /**
+ * @since 2.0.0
+ */
+ def withMaxConnectionAge(age: Duration): Http2ServerSettings =
copy(maxConnectionAge = age)
+
+ /**
+ * The time that requests in flight are given to complete after a connection
reached [[maxConnectionAge]]
+ * and the GOAWAY frame was sent. When the grace period expires, the
connection is closed even if requests
+ * are still in flight. The value `Duration.Inf` disables the limit, so that
the connection is closed only
+ * once all requests in flight have completed.
+ *
+ * @since 2.0.0
+ */
+ def maxConnectionAgeGrace: Duration
+
+ /**
+ * @since 2.0.0
+ */
+ def withMaxConnectionAgeGrace(grace: Duration): Http2ServerSettings =
copy(maxConnectionAgeGrace = grace)
+
+ /**
+ * The jitter applied to [[maxConnectionAge]] per connection, as a fraction
of the configured age: with
+ * the default of 0.1 each connection is closed after between 90% and 110%
of the configured age, so that
+ * connections that were opened together are not all closed at the same
time. 0 disables jitter.
+ *
+ * @since 2.0.0
+ */
+ def maxConnectionAgeJitter: Double
+
+ /**
+ * @since 2.0.0
+ */
+ override def withMaxConnectionAgeJitter(jitter: Double): Http2ServerSettings
=
+ copy(maxConnectionAgeJitter = jitter)
+
def frameTypeThrottleFrameTypes: Set[String]
def withFrameTypeThrottleFrameTypes(frameTypes: Set[String]) =
copy(frameTypeThrottleFrameTypes = frameTypes)
@@ -171,6 +218,9 @@ object Http2ServerSettings extends
SettingsCompanion[Http2ServerSettings] {
logFrames: Boolean,
pingInterval: FiniteDuration,
pingTimeout: FiniteDuration,
+ maxConnectionAge: Duration,
+ maxConnectionAgeGrace: Duration,
+ maxConnectionAgeJitter: Double,
frameTypeThrottleFrameTypes: Set[String],
frameTypeThrottleCost: Int,
frameTypeThrottleBurst: Int,
@@ -192,6 +242,10 @@ object Http2ServerSettings extends
SettingsCompanion[Http2ServerSettings] {
"min-collect-strict-entity-size <= incoming-connection-level-buffer-size
/ max-concurrent-streams")
require(outgoingControlFrameBufferSize > 0,
"outgoing-control-frame-buffer-size must be > 0")
require(frameTypeThrottleInterval.toMillis > 0,
"frame-type-throttle.interval must be a positive duration")
+ require(maxConnectionAge > Duration.Zero, "max-connection-age must be > 0
or 'infinite' to disable")
+ require(maxConnectionAgeGrace >= Duration.Zero, "max-connection-age-grace
must be >= 0 or 'infinite'")
+ require(maxConnectionAgeJitter >= 0 && maxConnectionAgeJitter < 1,
+ "max-connection-age-jitter must be >= 0 and < 1")
Http2CommonSettings.validate(this)
}
@@ -209,6 +263,9 @@ object Http2ServerSettings extends
SettingsCompanion[Http2ServerSettings] {
logFrames = c.getBoolean("log-frames"),
pingInterval = c.getFiniteDuration("ping-interval"),
pingTimeout = c.getFiniteDuration("ping-timeout"),
+ maxConnectionAge =
c.getPotentiallyInfiniteDuration("max-connection-age"),
+ maxConnectionAgeGrace =
c.getPotentiallyInfiniteDuration("max-connection-age-grace"),
+ maxConnectionAgeJitter = c.getDouble("max-connection-age-jitter"),
frameTypeThrottleFrameTypes =
c.getStringList("frame-type-throttle.frame-types").asScala.toSet,
frameTypeThrottleCost = c.getInt("frame-type-throttle.cost"),
frameTypeThrottleBurst = c.getInt("frame-type-throttle.burst"),
diff --git
a/http-core/src/test/scala/org/apache/pekko/http/scaladsl/settings/Http2ServerSettingsSpec.scala
b/http-core/src/test/scala/org/apache/pekko/http/scaladsl/settings/Http2ServerSettingsSpec.scala
new file mode 100644
index 000000000..9fed6f5d4
--- /dev/null
+++
b/http-core/src/test/scala/org/apache/pekko/http/scaladsl/settings/Http2ServerSettingsSpec.scala
@@ -0,0 +1,69 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.pekko.http.scaladsl.settings
+
+import java.time.temporal.ChronoUnit
+
+import org.apache.pekko.testkit.PekkoSpec
+
+import scala.concurrent.duration._
+
+class Http2ServerSettingsSpec extends PekkoSpec {
+
+ "Http2ServerSettings max-connection-age" should {
+
+ "be disabled by default" in {
+ val settings = Http2ServerSettings(system)
+ settings.maxConnectionAge should ===(Duration.Inf)
+ settings.getMaxConnectionAge should ===(ChronoUnit.FOREVER.getDuration)
+ }
+
+ "round-trip an infinite value through the Java API" in {
+ val settings =
Http2ServerSettings(system).withMaxConnectionAge(Duration.Inf)
+ settings.getMaxConnectionAge should ===(ChronoUnit.FOREVER.getDuration)
+ val roundTripped =
settings.withMaxConnectionAge(settings.getMaxConnectionAge)
+ roundTripped.getMaxConnectionAge should
===(ChronoUnit.FOREVER.getDuration)
+ }
+
+ "round-trip a finite value through the Java API" in {
+ val settings =
Http2ServerSettings(system).withMaxConnectionAge(2.minutes)
+ settings.getMaxConnectionAge should ===(java.time.Duration.ofMinutes(2))
+ val roundTripped =
settings.withMaxConnectionAge(settings.getMaxConnectionAge)
+ roundTripped.getMaxConnectionAge should
===(java.time.Duration.ofMinutes(2))
+ }
+ }
+
+ "Http2ServerSettings max-connection-age-grace" should {
+
+ "default to 30 seconds" in {
+ Http2ServerSettings(system).maxConnectionAgeGrace should ===(30.seconds)
+ }
+
+ "accept 'infinite' from config" in {
+ val settings =
Http2ServerSettings("pekko.http.server.http2.max-connection-age-grace =
infinite")
+ settings.maxConnectionAgeGrace should ===(Duration.Inf)
+ }
+
+ "round-trip an infinite value through the Java API" in {
+ val settings =
Http2ServerSettings(system).withMaxConnectionAgeGrace(Duration.Inf)
+ settings.getMaxConnectionAgeGrace should
===(ChronoUnit.FOREVER.getDuration)
+ val roundTripped =
settings.withMaxConnectionAgeGrace(settings.getMaxConnectionAgeGrace)
+ roundTripped.getMaxConnectionAgeGrace should
===(ChronoUnit.FOREVER.getDuration)
+ }
+ }
+}
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 2319ace06..2b6d19ab4 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
@@ -2113,6 +2113,125 @@ class Http2ServerSpec extends
Http2SpecWithMaterializer("""
terminated.futureValue
})
}
+ "support max-connection-age" should {
+ "send GOAWAY and close the connection when the age expires and no
requests are in flight".inAssertAllStagesStopped(
+ new TestSetup with RequestResponseProbes {
+ override def settings: ServerSettings = {
+ val default = super.settings
+ default.withHttp2Settings(
+
default.http2Settings.withMaxConnectionAge(500.millis).withMaxConnectionAgeJitter(0))
+ }
+
+ // with jitter disabled no GOAWAY is sent before the configured age:
checked for 400 of the 500 ms,
+ // the rest is left as a margin for timer scheduling
+ network.expectNoBytes(400.millis)
+ val (_, errorCode) = network.expectGOAWAY()
+ errorCode should ===(ErrorCode.NO_ERROR)
+ network.expectComplete()
+ })
+ "let requests in flight complete and refuse new streams when the age
expires".inAssertAllStagesStopped(
+ new TestSetup with RequestResponseProbes {
+ override def settings: ServerSettings = {
+ val default = super.settings
+
default.withHttp2Settings(default.http2Settings.withMaxConnectionAge(500.millis))
+ }
+
+ network.sendRequest(1, HttpRequest())
+ user.expectRequest()
+
+ val (_, errorCode) = network.expectGOAWAY(1)
+ errorCode should ===(ErrorCode.NO_ERROR)
+
+ // a stream opened after the GOAWAY was sent is refused, the client
is expected
+ // to retry it on a new connection
+ network.sendRequest(3, HttpRequest())
+ network.expectRST_STREAM(3, ErrorCode.REFUSED_STREAM)
+
+ // the request that was in flight when the age expired completes
normally
+ user.emitResponse(1, HttpResponse())
+ network.expectDecodedHEADERS(1)
+
+ network.expectComplete()
+ })
+ "close the connection when the grace period expires while requests are
still in flight".inAssertAllStagesStopped(
+ new TestSetup with RequestResponseProbes {
+ override def settings: ServerSettings = {
+ val default = super.settings
+ default.withHttp2Settings(
+
default.http2Settings.withMaxConnectionAge(500.millis).withMaxConnectionAgeGrace(300.millis))
+ }
+
+ network.sendRequest(1, HttpRequest())
+ user.expectRequest()
+
+ val (_, errorCode) = network.expectGOAWAY(1)
+ errorCode should ===(ErrorCode.NO_ERROR)
+
+ // the request in flight is never completed, the connection stays
open until the grace period expires
+ network.expectNoBytes(100.millis)
+ network.expectComplete()
+ })
+ "close within the deadline of a later server binding termination while
draining after the age expired".inAssertAllStagesStopped(
+ new TestSetup with RequestResponseProbes {
+ override def settings: ServerSettings = {
+ val default = super.settings
+ default.withHttp2Settings(
+
default.http2Settings.withMaxConnectionAge(500.millis).withMaxConnectionAgeJitter(0))
+ }
+
+ network.sendRequest(1, HttpRequest())
+ user.expectRequest()
+
+ val (_, errorCode) = network.expectGOAWAY(1)
+ errorCode should ===(ErrorCode.NO_ERROR)
+
+ // the default grace period is much longer than the deadline of the
termination, which must win
+ val terminated = serverTerminator.terminate(10.millis)
+ network.expectComplete()
+ terminated.futureValue
+ })
+ "keep the remaining grace period when a later server binding termination
has a later deadline".inAssertAllStagesStopped(
+ new TestSetup with RequestResponseProbes {
+ override def settings: ServerSettings = {
+ val default = super.settings
+ default.withHttp2Settings(
+
default.http2Settings.withMaxConnectionAge(500.millis).withMaxConnectionAgeJitter(0)
+ .withMaxConnectionAgeGrace(300.millis))
+ }
+
+ network.sendRequest(1, HttpRequest())
+ user.expectRequest()
+
+ val (_, errorCode) = network.expectGOAWAY(1)
+ errorCode should ===(ErrorCode.NO_ERROR)
+
+ val terminated = serverTerminator.terminate(1.minute)
+ network.expectNoBytes(100.millis)
+ network.expectComplete()
+ terminated.futureValue
+ })
+ "not shorten a server binding termination in progress when the age
expires".inAssertAllStagesStopped(
+ new TestSetup with RequestResponseProbes {
+ override def settings: ServerSettings = {
+ val default = super.settings
+ default.withHttp2Settings(
+
default.http2Settings.withMaxConnectionAge(300.millis).withMaxConnectionAgeJitter(0)
+ .withMaxConnectionAgeGrace(100.millis))
+ }
+
+ network.sendRequest(1, HttpRequest())
+ user.expectRequest()
+
+ val terminated = serverTerminator.terminate(1.second)
+ val (_, errorCode) = network.expectGOAWAY(1)
+ errorCode should ===(ErrorCode.NO_ERROR)
+
+ // the age expires while the termination is in progress, its grace
period must not apply
+ network.expectNoBytes(700.millis)
+ network.expectComplete()
+ terminated.futureValue
+ })
+ }
}
}
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]