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]

Reply via email to