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 561929f96 feat: add client-side max connection age for HTTP/2 (#1320)
561929f96 is described below

commit 561929f9659961d89ed7e6e5b97b701362c5a832
Author: Rayan-and-beyond <[email protected]>
AuthorDate: Sat Oct 3 13:58:36 2026 +0300

    feat: add client-side max connection age for HTTP/2 (#1320)
    
    * http2: add persistent client connection max age #1319
    
    Motivation:
    Long-lived managed HTTP/2 clients can stay pinned to existing server 
instances.
    
    Modification:
    Add a configurable maximum age that retires managed persistent HTTP/2 
connections after in-flight requests drain.
    
    Result:
    Requests arriving after retirement establish a fresh connection while 
in-flight requests complete normally.
    
    Tests:
    - sbt validatePullRequest
    - sbt http2-tests/test
    - sbt +http-core/mimaReportBinaryIssues
    - sbt scalafmtCheckAll scalafmtSbtCheck
    - sbt +headerCheckAll
    - sbt docs/paradox
    - sbt checkCodeStyle
    - git diff --check
    
    References:
    Refs #1319
    
    * http2: address max-age review feedback
    
    Expose client max-age and jitter as public settings, add per-connection 
jitter, and preserve a buffered request when the request source completes 
during retirement. Document break-before-make behavior and extend 
retirement/reconnect coverage.
    
    * http2: preserve Java max-age precision
    
    * Align client max connection age with the server-side setting from #1316
    
    Motivation:
    #1316 merged with `infinite` as the disabled value for
    `max-connection-age`, a `Duration` setting type and
    `JavaDurationConverter` for the Java API. The client setting used `0s`
    and `FiniteDuration`.
    
    Modification:
    - `persistent-connection-max-age` defaults to `infinite`, is a
      `Duration`, and must be > 0 or `infinite`, as on the server
    - Java accessors use `JavaDurationConverter`, so
      `ChronoUnit.FOREVER.getDuration` round-trips to `Duration.Inf`
    - reference.conf, scaladoc and docs follow the server wording
    - the jitter scheduling matches `Http2Demux`
    - settings tests move to a new `Http2ClientSettingsSpec`, mirroring
      `Http2ServerSettingsSpec`
    - MiMa excludes are reduced to the abstract members that need them
    - pekko-style imports in the changed files
    
    Result:
    The client and server max connection age settings share naming
    conventions, defaults, validation and Java conversion behaviour.
    
    Tests:
    - sbt "http-core/testOnly ...Http2ClientSettingsSpec 
...Http2CommonSettingsSpec ...Http2ServerSettingsSpec": 16 passed
    - sbt "http2-tests/testOnly ...Http2PersistentClient*": 26 passed
    - sbt "http-core/mimaReportBinaryIssues": clean
    - sbt scalafmtAll and headerCreateAll run on changed modules
    
    References:
    Refs #1319, #1316
    
    ---------
    
    Co-authored-by: PJ Fanning <[email protected]>
---
 docs/src/main/paradox/client-side/http2.md         | 29 +++++++
 .../http2-client-max-connection-age.excludes       | 21 +++++
 http-core/src/main/resources/reference.conf        | 23 ++++-
 .../pekko/http/impl/engine/http2/Http2Demux.scala  |  3 +-
 .../engine/http2/client/PersistentConnection.scala | 65 +++++++++++++--
 .../javadsl/settings/Http2ClientSettings.scala     | 38 ++++++++-
 .../scaladsl/settings/Http2ServerSettings.scala    | 40 +++++++++
 .../settings/Http2ClientSettingsSpec.scala         | 97 ++++++++++++++++++++++
 .../engine/http2/Http2PersistentClientSpec.scala   | 68 +++++++++++++++
 9 files changed, 372 insertions(+), 12 deletions(-)

diff --git a/docs/src/main/paradox/client-side/http2.md 
b/docs/src/main/paradox/client-side/http2.md
index 69261215e..125844621 100644
--- a/docs/src/main/paradox/client-side/http2.md
+++ b/docs/src/main/paradox/client-side/http2.md
@@ -55,6 +55,35 @@ Java
 
 The Apache Pekko HTTP client doesn't support HTTP/1 to HTTP/2 negotiation over 
plaintext using the `Upgrade` mechanism.
 
+## Limiting managed persistent connection lifetime
+
+Managed persistent HTTP/2 clients can periodically retire long-lived 
connections so that later requests establish
+fresh connections. This is useful when clients would otherwise remain pinned 
to the same server instances after
+a scale-out or rolling deployment.
+
+Configure the maximum age and optional jitter under the HTTP/2 client settings:
+
+```
+pekko.http.client.http2.persistent-connection-max-age = 10m
+pekko.http.client.http2.persistent-connection-max-age-jitter = 0.1
+```
+
+The setting applies to `managedPersistentHttp2()` and 
`managedPersistentHttp2WithPriorKnowledge()`. Each
+connection gets an independently jittered age; with the default jitter of 
`0.1`, retirement happens between 90%
+and 110% of the configured maximum age. Set the jitter to `0` to disable it.
+
+Retirement is currently **break-before-make**. When a connection reaches its 
age, the managed client stops
+assigning new requests to it and lets requests already in flight drain before 
closing the connection. New requests
+are backpressured until the old connection disconnects, then a replacement 
connection is established. As a result,
+a retiring connection can add its drain time plus TCP/TLS/HTTP2 setup time to 
new-request latency.
+
+The existing `completion-timeout` bounds the drain period. If that timeout 
expires, the old connection is closed
+even when requests are still in flight, terminating long-lived requests such 
as streaming responses.
+
+The default maximum age is `infinite`, which disables age-based retirement. 
The settings can also be changed
+programmatically with `Http2ClientSettings.withPersistentConnectionMaxAge` and
+`Http2ClientSettings.withPersistentConnectionMaxAgeJitter`.
+
 ## Request-response ordering
 
 For HTTP/2 connections the responses are not guaranteed to arrive in the same 
order that the requests were emitted to
diff --git 
a/http-core/src/main/mima-filters/2.0.x.backwards.excludes/http2-client-max-connection-age.excludes
 
b/http-core/src/main/mima-filters/2.0.x.backwards.excludes/http2-client-max-connection-age.excludes
new file mode 100644
index 000000000..be4f838ef
--- /dev/null
+++ 
b/http-core/src/main/mima-filters/2.0.x.backwards.excludes/http2-client-max-connection-age.excludes
@@ -0,0 +1,21 @@
+# 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 persistent connection max-age settings for HTTP/2 clients
+ProblemFilters.exclude[ReversedMissingMethodProblem]("org.apache.pekko.http.scaladsl.settings.Http2ClientSettings.persistentConnectionMaxAge")
+ProblemFilters.exclude[ReversedMissingMethodProblem]("org.apache.pekko.http.scaladsl.settings.Http2ClientSettings.persistentConnectionMaxAgeJitter")
+ProblemFilters.exclude[ReversedMissingMethodProblem]("org.apache.pekko.http.javadsl.settings.Http2ClientSettings.withPersistentConnectionMaxAgeJitter")
diff --git a/http-core/src/main/resources/reference.conf 
b/http-core/src/main/resources/reference.conf
index e591cead0..7fe20b171 100644
--- a/http-core/src/main/resources/reference.conf
+++ b/http-core/src/main/resources/reference.conf
@@ -633,6 +633,26 @@ pekko.http {
       # Set to zero to retry indefinitely.
       max-persistent-attempts = 0
 
+      # The maximum time a connection created by managedPersistentHttp2 or 
managedPersistentHttp2WithPriorKnowledge
+      # is used before it is retired. When the age of a connection exceeds 
this value, the connection stops
+      # accepting new requests, lets requests that are already in flight 
complete (see `completion-timeout`), and
+      # is then closed. Later requests are sent on a new connection.
+      #
+      # Retiring 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.
+      #
+      # Retirement is break-before-make: new requests wait for the retiring 
connection to close before a
+      # new connection is established, so they can observe the drain plus the 
connection setup latency.
+      #
+      # The value `infinite` disables this mechanism and is the default.
+      persistent-connection-max-age = infinite
+
+      # The jitter applied to `persistent-connection-max-age`, as a fraction 
of the configured value: with the
+      # default of 0.1 each connection is retired after between 90% and 110% 
of the configured age, so that
+      # connections that were opened together are not all retired at the same 
time.
+      # Set to 0 to disable jitter. Must be >= 0 and < 1.
+      persistent-connection-max-age-jitter = 0.1
+
       # Starting backoff before reconnecting when a persistent HTTP/2 client 
connection fails
       # see `pekko.http.host-connection-pool.base-connection-backoff` for 
details on the backoff mechanism.
       base-connection-backoff = 
${pekko.http.host-connection-pool.base-connection-backoff}
@@ -642,7 +662,8 @@ pekko.http {
       max-connection-backoff = 
${pekko.http.host-connection-pool.max-connection-backoff}
 
       # When gracefully closing the HTTP/2 client, await at most 
`completion-timeout` for in-flight
-      # requests to complete.
+      # requests to complete. When the timeout expires, remaining in-flight 
requests are terminated.
+      # This also bounds the drain of a connection retired after 
`persistent-connection-max-age`.
       completion-timeout = 3s
 
     }
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 4e3cd4854..cee6b45d6 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
@@ -72,7 +72,8 @@ private[http2] class Http2ClientDemux(http2Settings: 
Http2ClientSettings, master
 
   override def completionTimeout: FiniteDuration = 
http2Settings.completionTimeout
 
-  // a maximum connection age is not supported on the client side
+  // the client side limits the age of managed persistent connections in 
PersistentConnection instead,
+  // see `persistent-connection-max-age`
   def maxConnectionAge: Duration = Duration.Inf
   def maxConnectionAgeGrace: Duration = Duration.Inf
   def maxConnectionAgeJitter: Double = 0.0
diff --git 
a/http-core/src/main/scala/org/apache/pekko/http/impl/engine/http2/client/PersistentConnection.scala
 
b/http-core/src/main/scala/org/apache/pekko/http/impl/engine/http2/client/PersistentConnection.scala
index 6f8f8db39..1731ab6e8 100644
--- 
a/http-core/src/main/scala/org/apache/pekko/http/impl/engine/http2/client/PersistentConnection.scala
+++ 
b/http-core/src/main/scala/org/apache/pekko/http/impl/engine/http2/client/PersistentConnection.scala
@@ -35,6 +35,7 @@ import scala.util.{ Failure, Success }
 private[http2] object PersistentConnection {
 
   private case class EmbargoEnded(connectsLeft: Option[Int], embargo: 
FiniteDuration)
+  private case object MaxConnectionAgeReached
 
   /**
    * Wraps a connection flow with transparent reconnection support.
@@ -59,7 +60,8 @@ private[http2] object PersistentConnection {
       settings.maxPersistentAttempts match {
         case 0 => None
         case n => Some(n)
-      }, settings.baseConnectionBackoff, settings.maxConnectionBackoff))
+      }, settings.baseConnectionBackoff, settings.maxConnectionBackoff,
+      settings.persistentConnectionMaxAge, 
settings.persistentConnectionMaxAgeJitter))
 
   private class AssociationTag extends RequestResponseAssociation
   private val associationTagKey = 
AttributeKey[AssociationTag]("PersistentConnection.associationTagKey")
@@ -69,7 +71,8 @@ private[http2] object PersistentConnection {
       entity = "The server closed the connection before delivering a 
response.")
 
   private class Stage(connectionFlow: Flow[HttpRequest, HttpResponse, 
Future[OutgoingConnection]],
-      maxAttempts: Option[Int], baseEmbargo: FiniteDuration, _maxBackoff: 
FiniteDuration)
+      maxAttempts: Option[Int], baseEmbargo: FiniteDuration, _maxBackoff: 
FiniteDuration,
+      persistentConnectionMaxAge: Duration, persistentConnectionMaxAgeJitter: 
Double)
       extends GraphStage[FlowShape[HttpRequest, HttpResponse]] {
     val requestIn = Inlet[HttpRequest]("PersistentConnection.requestIn")
     val responseOut = Outlet[HttpResponse]("PersistentConnection.responseOut")
@@ -78,6 +81,8 @@ private[http2] object PersistentConnection {
     val shape: FlowShape[HttpRequest, HttpResponse] = FlowShape(requestIn, 
responseOut)
     override def createLogic(inheritedAttributes: Attributes): GraphStageLogic 
=
       new TimerGraphStageLogic(shape) with StageLogging {
+        private var currentConnection: Option[Connected] = None
+
         become(Unconnected)
 
         def become(state: State): Unit = setHandlers(requestIn, responseOut, 
state)
@@ -86,7 +91,9 @@ private[http2] object PersistentConnection {
         object Unconnected extends State {
           override def onPush(): Unit = connect(maxAttempts, Duration.Zero)
           override def onPull(): Unit =
-            if (!isAvailable(requestIn) && !hasBeenPulled(requestIn)) // 
requestIn might already have been pulled when we failed and went back to 
Unconnected
+            if (isAvailable(requestIn)) connect(maxAttempts, Duration.Zero)
+            else if (isClosed(requestIn)) completeStage()
+            else if (!hasBeenPulled(requestIn)) // requestIn might already 
have been pulled when we failed and went back to Unconnected
               pull(requestIn)
         }
 
@@ -131,12 +138,13 @@ private[http2] object PersistentConnection {
           override def onPush(): Unit = () // Pull might have happened before 
the connection failed. Element is kept in slot.
 
           override def onPull(): Unit = {
-            if (!isAvailable(requestIn) && !hasBeenPulled(requestIn)) // 
requestIn might already have been pulled when we failed and went back to 
Unconnected
+            if (!isAvailable(requestIn) && !isClosed(requestIn) && 
!hasBeenPulled(requestIn)) // requestIn might already have been pulled when we 
failed and went back to Unconnected
               pull(requestIn)
           }
 
           val onConnected = getAsyncCallback[Unit] { _ =>
             val newState = new Connected(requestOut, responseIn)
+            currentConnection = Some(newState)
             become(newState)
             if (requestOutPulled) {
               if (isAvailable(requestIn)) 
newState.dispatchRequest(grab(requestIn))
@@ -181,6 +189,8 @@ private[http2] object PersistentConnection {
             case EmbargoEnded(connectsLeft, nextEmbargo) =>
               log.debug("Reconnecting after backoff")
               connect(connectsLeft, nextEmbargo)
+            case MaxConnectionAgeReached =>
+              currentConnection.foreach(_.retire())
           }
         }
 
@@ -188,12 +198,27 @@ private[http2] object PersistentConnection {
             requestOut: SubSourceOutlet[HttpRequest],
             responseIn: SubSinkInlet[HttpResponse]) extends State {
           private var ongoingRequests: Map[AssociationTag, 
Map[AttributeKey[?], RequestResponseAssociation]] = Map.empty
+          private var retiring = false
+
+          persistentConnectionMaxAge match {
+            case age: FiniteDuration =>
+              // The age of each connection is jittered so that connections 
that were opened together are not
+              // all retired at the same time, see 
`persistent-connection-max-age-jitter` in the configuration.
+              val jitterFactor =
+                1.0 + persistentConnectionMaxAgeJitter * (2 * 
ThreadLocalRandom.current().nextDouble() - 1)
+              scheduleOnce(MaxConnectionAgeReached, (age.toMillis * 
jitterFactor).toLong.max(1L).millis)
+            case _ => // no maximum connection age configured
+          }
+
           responseIn.pull()
 
           requestOut.setHandler(new OutHandler {
             override def onPull(): Unit =
-              if (!isAvailable(requestIn)) pull(requestIn)
-              else dispatchRequest(grab(requestIn))
+              if (isAvailable(requestIn)) {
+                dispatchRequest(grab(requestIn))
+                if (isClosed(requestIn)) requestOut.complete()
+              } else if (isClosed(requestIn)) requestOut.complete()
+              else if (!hasBeenPulled(requestIn)) pull(requestIn)
 
             override def onDownstreamFinish(cause: Throwable): Unit = 
onDisconnected()
           })
@@ -210,23 +235,45 @@ private[http2] object PersistentConnection {
             override def onUpstreamFailure(ex: Throwable): Unit = 
onDisconnected() // FIXME: log error
           })
           def onDisconnected(): Unit = {
+            cancelTimer(MaxConnectionAgeReached)
+            currentConnection = None
+
             emitMultiple[HttpResponse](responseOut,
               
ongoingRequests.values.map(errorResponse.withAttributes(_)).toVector,
               () => setHandler(responseOut, Unconnected))
             responseIn.cancel()
             requestOut.fail(new RuntimeException("connection broken"))
 
-            if (isClosed(requestIn)) {
-              // user closed PersistentConnection before and we were waiting 
for remaining responses
+            if (isClosed(requestIn) && !isAvailable(requestIn)) {
+              // user closed PersistentConnection before and there is no 
buffered request left to dispatch
               completeStage()
             } else {
               // become(Unconnected) doesn't work because of using emit
               // so we need to do it more carefully here
               setHandler(requestIn, Unconnected)
-              if (isAvailable(responseOut) && !hasBeenPulled(requestIn)) 
pull(requestIn)
+              if (isAvailable(responseOut)) {
+                if (isAvailable(requestIn)) connect(maxAttempts, Duration.Zero)
+                else if (!hasBeenPulled(requestIn)) pull(requestIn)
+              }
             }
           }
 
+          def retire(): Unit =
+            if (!retiring) {
+              retiring = true
+              log.debug("Persistent HTTP/2 connection reached its configured 
maximum age, retiring it")
+              setHandler(requestIn,
+                new InHandler {
+                  override def onPush(): Unit = () // keep at most one next 
request in the inlet slot
+                  override def onUpstreamFinish(): Unit = ()
+                  override def onUpstreamFailure(ex: Throwable): Unit = {
+                    responseIn.cancel()
+                    failStage(ex)
+                  }
+                })
+              requestOut.complete()
+            }
+
           def dispatchRequest(req: HttpRequest): Unit = {
             val tag = new AssociationTag
             // Some cross-compilation woes here:
diff --git 
a/http-core/src/main/scala/org/apache/pekko/http/javadsl/settings/Http2ClientSettings.scala
 
b/http-core/src/main/scala/org/apache/pekko/http/javadsl/settings/Http2ClientSettings.scala
index 230413739..f7e6e5fec 100644
--- 
a/http-core/src/main/scala/org/apache/pekko/http/javadsl/settings/Http2ClientSettings.scala
+++ 
b/http-core/src/main/scala/org/apache/pekko/http/javadsl/settings/Http2ClientSettings.scala
@@ -15,7 +15,9 @@ package org.apache.pekko.http.javadsl.settings
 
 import java.time.Duration
 
-import org.apache.pekko.http.scaladsl
+import org.apache.pekko
+import pekko.http.impl.util.JavaDurationConverter
+import pekko.http.scaladsl
 
 import scala.concurrent.duration.DurationLong
 
@@ -79,6 +81,40 @@ trait Http2ClientSettings { self: 
scaladsl.settings.Http2ClientSettings.Http2Cli
   def getMaxPersistentAttempts: Int = maxPersistentAttempts
   def withMaxPersistentAttempts(max: Int): Http2ClientSettings = 
copy(maxPersistentAttempts = max)
 
+  /**
+   * The maximum age of a connection created by `managedPersistentHttp2` or
+   * `managedPersistentHttp2WithPriorKnowledge`. When the age of a connection 
exceeds this value, the connection
+   * stops accepting new requests, lets requests that are already in flight 
complete within
+   * [[getCompletionTimeout]], and is then closed. Later requests are sent on 
a new connection. The value
+   * `ChronoUnit.FOREVER.getDuration` represents an infinite age, which 
disables this mechanism and is the
+   * default.
+   *
+   * @since 2.0.0
+   */
+  def getPersistentConnectionMaxAge: Duration = 
JavaDurationConverter.toJava(persistentConnectionMaxAge)
+
+  /**
+   * Pass `ChronoUnit.FOREVER.getDuration` to disable the maximum connection 
age.
+   *
+   * @since 2.0.0
+   */
+  def withPersistentConnectionMaxAge(maxAge: Duration): Http2ClientSettings =
+    self.withPersistentConnectionMaxAge(JavaDurationConverter.toScala(maxAge))
+
+  /**
+   * The jitter applied to the persistent connection maximum age per 
connection, as a fraction of the configured
+   * age: with the default of 0.1 each connection is retired after between 90% 
and 110% of the configured age, so
+   * that connections that were opened together are not all retired at the 
same time. 0 disables jitter.
+   *
+   * @since 2.0.0
+   */
+  def getPersistentConnectionMaxAgeJitter: Double = 
persistentConnectionMaxAgeJitter
+
+  /**
+   * @since 2.0.0
+   */
+  def withPersistentConnectionMaxAgeJitter(jitter: Double): Http2ClientSettings
+
   def getCompletionTimeout: Duration = 
Duration.ofMillis(completionTimeout.toMillis)
   def withCompletionTimeout(timeout: Duration): Http2ClientSettings = 
copy(completionTimeout = timeout.toMillis.millis)
 
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 7d5b42bd3..e2616650d 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
@@ -347,6 +347,38 @@ trait Http2ClientSettings extends 
javadsl.settings.Http2ClientSettings with Http
   def maxPersistentAttempts: Int
   override def withMaxPersistentAttempts(max: Int): Http2ClientSettings = 
copy(maxPersistentAttempts = max)
 
+  /**
+   * The maximum age of a connection created by `managedPersistentHttp2` or
+   * `managedPersistentHttp2WithPriorKnowledge`. When the age of a connection 
exceeds this value, the connection
+   * stops accepting new requests, lets requests that are already in flight 
complete within [[completionTimeout]],
+   * and is then closed. Later requests are sent on a new connection. The 
value `Duration.Inf` disables this
+   * mechanism and is the default.
+   *
+   * @since 2.0.0
+   */
+  def persistentConnectionMaxAge: Duration
+
+  /**
+   * @since 2.0.0
+   */
+  def withPersistentConnectionMaxAge(maxAge: Duration): Http2ClientSettings =
+    copy(persistentConnectionMaxAge = maxAge)
+
+  /**
+   * The jitter applied to [[persistentConnectionMaxAge]] per connection, as a 
fraction of the configured age:
+   * with the default of 0.1 each connection is retired after between 90% and 
110% of the configured age, so
+   * that connections that were opened together are not all retired at the 
same time. 0 disables jitter.
+   *
+   * @since 2.0.0
+   */
+  def persistentConnectionMaxAgeJitter: Double
+
+  /**
+   * @since 2.0.0
+   */
+  override def withPersistentConnectionMaxAgeJitter(jitter: Double): 
Http2ClientSettings =
+    copy(persistentConnectionMaxAgeJitter = jitter)
+
   def completionTimeout: FiniteDuration
   def withCompletionTimeout(timeout: FiniteDuration): Http2ClientSettings = 
copy(completionTimeout = timeout)
 
@@ -380,6 +412,8 @@ object Http2ClientSettings extends 
SettingsCompanion[Http2ClientSettings] {
       pingInterval: FiniteDuration,
       pingTimeout: FiniteDuration,
       maxPersistentAttempts: Int,
+      persistentConnectionMaxAge: Duration,
+      persistentConnectionMaxAgeJitter: Double,
       completionTimeout: FiniteDuration,
       baseConnectionBackoff: FiniteDuration,
       maxConnectionBackoff: FiniteDuration,
@@ -395,6 +429,10 @@ object Http2ClientSettings extends 
SettingsCompanion[Http2ClientSettings] {
     require(incomingStreamLevelBufferSize > 0, 
"incoming-stream-level-buffer-size must be > 0")
     require(outgoingControlFrameBufferSize > 0, 
"outgoing-control-frame-buffer-size must be > 0")
     require(maxPersistentAttempts >= 0, "max-persistent-attempts must be >= 0")
+    require(persistentConnectionMaxAge > Duration.Zero,
+      "persistent-connection-max-age must be > 0 or 'infinite' to disable")
+    require(persistentConnectionMaxAgeJitter >= 0 && 
persistentConnectionMaxAgeJitter < 1,
+      "persistent-connection-max-age-jitter must be >= 0 and < 1")
     require(completionTimeout > Duration.Zero, "completion-timeout must be > 
0")
     require(baseConnectionBackoff <= maxConnectionBackoff, 
"base-connection-backoff must be <= max-connection-backoff")
     Http2CommonSettings.validate(this)
@@ -414,6 +452,8 @@ object Http2ClientSettings extends 
SettingsCompanion[Http2ClientSettings] {
       pingInterval = c.getFiniteDuration("ping-interval"),
       pingTimeout = c.getFiniteDuration("ping-timeout"),
       maxPersistentAttempts = c.getInt("max-persistent-attempts"),
+      persistentConnectionMaxAge = 
c.getPotentiallyInfiniteDuration("persistent-connection-max-age"),
+      persistentConnectionMaxAgeJitter = 
c.getDouble("persistent-connection-max-age-jitter"),
       completionTimeout = c.getFiniteDuration("completion-timeout"),
       baseConnectionBackoff = c.getFiniteDuration("base-connection-backoff"),
       maxConnectionBackoff = c.getFiniteDuration("max-connection-backoff"),
diff --git 
a/http-core/src/test/scala/org/apache/pekko/http/scaladsl/settings/Http2ClientSettingsSpec.scala
 
b/http-core/src/test/scala/org/apache/pekko/http/scaladsl/settings/Http2ClientSettingsSpec.scala
new file mode 100644
index 000000000..e4222e028
--- /dev/null
+++ 
b/http-core/src/test/scala/org/apache/pekko/http/scaladsl/settings/Http2ClientSettingsSpec.scala
@@ -0,0 +1,97 @@
+/*
+ * 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
+import pekko.testkit.PekkoSpec
+
+import scala.concurrent.duration._
+
+class Http2ClientSettingsSpec extends PekkoSpec {
+
+  "Http2ClientSettings persistent-connection-max-age" should {
+
+    "be disabled by default" in {
+      val settings = Http2ClientSettings(system)
+      settings.persistentConnectionMaxAge should ===(Duration.Inf)
+      settings.getPersistentConnectionMaxAge should 
===(ChronoUnit.FOREVER.getDuration)
+      settings.internalSettings should ===(None)
+    }
+
+    "accept a finite value from config" in {
+      val settings = 
Http2ClientSettings("pekko.http.client.http2.persistent-connection-max-age = 
2s")
+      settings.persistentConnectionMaxAge should ===(2.seconds)
+      settings.getPersistentConnectionMaxAge should 
===(java.time.Duration.ofSeconds(2))
+    }
+
+    "round-trip an infinite value through the Java API" in {
+      val settings = 
Http2ClientSettings(system).withPersistentConnectionMaxAge(Duration.Inf)
+      settings.getPersistentConnectionMaxAge should 
===(ChronoUnit.FOREVER.getDuration)
+      val roundTripped = 
settings.withPersistentConnectionMaxAge(settings.getPersistentConnectionMaxAge)
+      roundTripped.getPersistentConnectionMaxAge should 
===(ChronoUnit.FOREVER.getDuration)
+    }
+
+    "round-trip a finite value through the Java API" in {
+      val settings = 
Http2ClientSettings(system).withPersistentConnectionMaxAge(2.minutes)
+      settings.getPersistentConnectionMaxAge should 
===(java.time.Duration.ofMinutes(2))
+      val roundTripped = 
settings.withPersistentConnectionMaxAge(settings.getPersistentConnectionMaxAge)
+      roundTripped.getPersistentConnectionMaxAge should 
===(java.time.Duration.ofMinutes(2))
+    }
+
+    "keep sub-millisecond precision through the Java API" in {
+      val age = java.time.Duration.ofNanos(1)
+      
Http2ClientSettings(system).withPersistentConnectionMaxAge(age).getPersistentConnectionMaxAge
 should ===(age)
+    }
+
+    "reject a zero or negative value" in {
+      intercept[IllegalArgumentException] {
+        
Http2ClientSettings("pekko.http.client.http2.persistent-connection-max-age = 
0s")
+      }
+      intercept[IllegalArgumentException] {
+        
Http2ClientSettings("pekko.http.client.http2.persistent-connection-max-age = 
-1s")
+      }
+    }
+  }
+
+  "Http2ClientSettings persistent-connection-max-age-jitter" should {
+
+    "default to 0.1" in {
+      val settings = Http2ClientSettings(system)
+      settings.persistentConnectionMaxAgeJitter should ===(0.1)
+      settings.getPersistentConnectionMaxAgeJitter should ===(0.1)
+    }
+
+    "accept a value from config and programmatically" in {
+      
Http2ClientSettings("pekko.http.client.http2.persistent-connection-max-age-jitter
 = 0.25")
+        .persistentConnectionMaxAgeJitter should ===(0.25)
+      Http2ClientSettings(system).withPersistentConnectionMaxAgeJitter(0.2)
+        .persistentConnectionMaxAgeJitter should ===(0.2)
+    }
+
+    "reject a value outside of [0, 1)" in {
+      intercept[IllegalArgumentException] {
+        
Http2ClientSettings("pekko.http.client.http2.persistent-connection-max-age-jitter
 = -0.1")
+      }
+      intercept[IllegalArgumentException] {
+        
Http2ClientSettings("pekko.http.client.http2.persistent-connection-max-age-jitter
 = 1")
+      }
+    }
+  }
+}
diff --git 
a/http2-tests/src/test/scala/org/apache/pekko/http/impl/engine/http2/Http2PersistentClientSpec.scala
 
b/http2-tests/src/test/scala/org/apache/pekko/http/impl/engine/http2/Http2PersistentClientSpec.scala
index 7378681de..14bb3899f 100644
--- 
a/http2-tests/src/test/scala/org/apache/pekko/http/impl/engine/http2/Http2PersistentClientSpec.scala
+++ 
b/http2-tests/src/test/scala/org/apache/pekko/http/impl/engine/http2/Http2PersistentClientSpec.scala
@@ -102,6 +102,73 @@ abstract class Http2PersistentClientSpec(tls: Boolean) 
extends PekkoSpecWithMate
         response.attribute(requestIdAttr).get.id shouldBe "request-1"
       })
 
+    "retire an aged connection after in-flight requests 
complete".inAssertAllStagesStopped(new TestSetup(tls) {
+      override def clientSettings: ClientConnectionSettings =
+        super.clientSettings.withHttp2Settings(Http2ClientSettings(
+          """
+            pekko.http.client.http2.persistent-connection-max-age = 500ms
+            pekko.http.client.http2.persistent-connection-max-age-jitter = 0
+            pekko.http.client.http2.completion-timeout = 2s
+          """))
+
+      client.responsesIn.request(2)
+      client.sendRequest(HttpRequest(uri = 
"/first").addAttribute(requestIdAttr, RequestId("request-1")))
+
+      val first = server.expectRequest()
+      val firstClientPort = first.clientPort
+      killProbe.expectMsgType[UniqueKillSwitch]
+
+      // Let the configured age elapse while the first request is still in 
flight. Use a wide margin so this
+      // remains stable on slow CI workers.
+      server.expectNoRequest(1500.millis)
+
+      client.sendRequest(HttpRequest(uri = 
"/second").addAttribute(requestIdAttr, RequestId("request-2")))
+      // Closing the request source with a buffered request must not drop that 
request during retirement.
+      client.requestsOut.sendComplete()
+      // A request arriving after retirement starts must wait for the current 
in-flight request.
+      server.expectNoRequest(300.millis)
+
+      server.sendResponseFor(first, HttpResponse(entity = "first-response"))
+      val firstResponse = client.expectResponse()
+      Unmarshal(firstResponse.entity).to[String].futureValue shouldBe 
"first-response"
+      firstResponse.attribute(requestIdAttr).get.id shouldBe "request-1"
+
+      val second = server.expectRequest()
+      second.clientPort should not be firstClientPort
+      server.sendResponseFor(second, HttpResponse(entity = "second-response"))
+
+      val secondResponse = client.expectResponse()
+      Unmarshal(secondResponse.entity).to[String].futureValue shouldBe 
"second-response"
+      secondResponse.attribute(requestIdAttr).get.id shouldBe "request-2"
+    })
+
+    "retire an idle aged connection and reconnect lazily for the next 
request".inAssertAllStagesStopped(
+      new TestSetup(tls) {
+        override def clientSettings: ClientConnectionSettings =
+          super.clientSettings.withHttp2Settings(Http2ClientSettings(
+            """
+              pekko.http.client.http2.persistent-connection-max-age = 500ms
+              pekko.http.client.http2.persistent-connection-max-age-jitter = 0
+            """))
+
+        client.responsesIn.request(2)
+        client.sendRequest(HttpRequest(uri = "/first"))
+        val first = server.expectRequest()
+        val firstClientPort = first.clientPort
+        server.sendResponseFor(first, HttpResponse())
+        client.expectResponse()
+
+        // Allow the idle connection to reach its configured age with a 
generous scheduling margin.
+        server.expectNoRequest(1500.millis)
+
+        client.sendRequest(HttpRequest(uri = "/second"))
+        val second = server.expectRequest()
+        second.clientPort should not be firstClientPort
+        server.sendResponseFor(second, HttpResponse())
+        client.expectResponse()
+        client.requestsOut.sendComplete()
+      })
+
     def reconnectionTests(withBackoff: Boolean): Unit = {
       val changeSettings: Http2ClientSettings => Http2ClientSettings =
         if (withBackoff) s => 
s.withBaseConnectionBackoff(300.millis).withMaxConnectionBackoff(800.millis)
@@ -396,6 +463,7 @@ abstract class Http2PersistentClientSpec(tls: Boolean) 
extends PekkoSpecWithMate
       }
 
       def expectRequest(): ServerRequest = 
requestProbe.expectMsgType[ServerRequest]
+      def expectNoRequest(duration: FiniteDuration): Unit = 
requestProbe.expectNoMessage(duration)
       def sendResponseFor(request: ServerRequest, response: HttpResponse): 
Unit =
         request.sendResponse(response)
     }


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to